环球热门:Go+Kafka实现延迟消息的实现示例
目录
前言原理简单的实现生产者延迟服务消费者改进点通用的延迟服务生产者负责延迟服务总结前言
延迟队列是一个非常有用的工具,我们经常遇到需要使用延迟队列的场景,比如延迟通知,订单关闭等等。
【资料图】
这篇文章主要是使用Go+Kafka实现延迟消息。
使用了sarama客户端。
原理
Kafka实现延迟消息分为下面三步:
生产者把消息发送到延迟队列延迟服务把延迟队列里超过延迟时间的消息写入真实队列消费者消费真实队列里的消息
简单的实现
生产者
生产者只是把消息发送到延迟队列
msg := &sarama.ProducerMessage{
Topic: kafka_delay_queue_test.DelayTopic,
Value: sarama.ByteEncoder("test" + strconv.Itoa(i)),
}
if _, _, err := producer.SendMessage(msg); err != nil {
log.Println(err)
}延迟服务
延迟服务会订阅延迟队列的消息,并把超时消息发送到真实队列
if err = consumerGroup.Consume(context.Background(),
[]string{kafka_delay_queue_test.DelayTopic}, consumer); err != nil {
break
}type Consumer struct {
producer sarama.SyncProducer
delay time.Duration
}
func NewConsumer(producer sarama.SyncProducer, delay time.Duration) *Consumer {
return &Consumer{
producer: producer,
delay: delay,
}
}
func (c *Consumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for message := range claim.Messages() {
// 如果消息已经超时,把消息发送到真实队列
now := time.Now()
if now.Sub(message.Timestamp) >= c.delay {
_, _, err := c.producer.SendMessage(&sarama.ProducerMessage{
Topic: kafka_delay_queue_test.RealTopic,
Key: sarama.ByteEncoder(message.Key),
Value: sarama.ByteEncoder(message.Value),
})
if err == nil {
session.MarkMessage(message, "")
}
continue
}
// 否则休眠一秒
time.Sleep(time.Second)
return nil
}
return nil
}消费者
消费者只是订阅真实队列并消费消息
if err = consumerGroup.Consume(context.Background(),
[]string{kafka_delay_queue_test.RealTopic}, consumer); err != nil {
break
}type Consumer struct{}
func NewConsumer() *Consumer {
return &Consumer{}
}
func (c *Consumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for message := range claim.Messages() {
fmt.Println("收到消息:", message.Value, message.Timestamp)
session.MarkMessage(message, "")
}
return nil
}改进点
通用的延迟服务
可以把延迟服务封装成一个通用的服务,这样生产者可以直接把消息发送给延迟服务,让延迟服务去处理剩下的逻辑。
延迟服务可以提供多个延时等级,比如5s、10s、30s、1m、5m、10m、1h、2h等,类似于RocketMQ。
生产者负责延迟服务
也可以让生产者负责延迟服务,让生产者自己把延迟队列里面的消息发送到真实队列。
下面是一个简单的实现:
// KafkaDelayQueueProducer 延迟队列生产者,包含了生产者和延迟服务
type KafkaDelayQueueProducer struct {
producer sarama.SyncProducer // 生产者
delayTopic string // 延迟服务主题
}
// NewKafkaDelayQueueProducer 创建延迟队列生产者
// producer 生产者
// delayServiceConsumerGroup 延迟服务消费者
// delayTime 延迟时间
// delayTopic 延迟服务主题
// realTopic 真实队列主题
func NewKafkaDelayQueueProducer(producer sarama.SyncProducer, delayServiceConsumerGroup sarama.ConsumerGroup,
delayTime time.Duration, delayTopic, realTopic string) *KafkaDelayQueueProducer {
// 启动延迟服务
consumer := NewDelayServiceConsumer(producer, delayTime, realTopic)
go func() {
for {
if err := delayServiceConsumerGroup.Consume(context.Background(),
[]string{delayTopic}, consumer); err != nil {
break
}
}
}()
return &KafkaDelayQueueProducer{
producer: producer,
delayTopic: delayTopic,
}
}
// SendMessage 发送消息
func (q *KafkaDelayQueueProducer) SendMessage(msg *sarama.ProducerMessage) (partition int32, offset int64, err error) {
msg.Topic = q.delayTopic
return q.producer.SendMessage(msg)
}
// DelayServiceConsumer 延迟服务消费者
type DelayServiceConsumer struct {
producer sarama.SyncProducer
delay time.Duration
realTopic string
}
func NewDelayServiceConsumer(producer sarama.SyncProducer, delay time.Duration,
realTopic string) *DelayServiceConsumer {
return &DelayServiceConsumer{
producer: producer,
delay: delay,
realTopic: realTopic,
}
}
func (c *DelayServiceConsumer) ConsumeClaim(session sarama.ConsumerGroupSession,
claim sarama.ConsumerGroupClaim) error {
for message := range claim.Messages() {
// 如果消息已经超时,把消息发送到真实队列
now := time.Now()
if now.Sub(message.Timestamp) >= c.delay {
_, _, err := c.producer.SendMessage(&sarama.ProducerMessage{
Topic: c.realTopic,
Key: sarama.ByteEncoder(message.Key),
Value: sarama.ByteEncoder(message.Value),
})
if err == nil {
session.MarkMessage(message, "")
}
continue
}
// 否则休眠一秒
time.Sleep(time.Second)
return nil
}
return nil
}
func (c *DelayServiceConsumer) Setup(sarama.ConsumerGroupSession) error {
return nil
}
func (c *DelayServiceConsumer) Cleanup(sarama.ConsumerGroupSession) error {
return nil
}总结
使用中间队列+轮询可以很容易的在Kafka实现延迟消息,如果需要一个通用的延迟队列也可以实现一个通用的延迟服务,也可以让消费者负责延迟服务的功能。
完整代码:
简单实现例子:https://github.com/jiaxwu/dq/blob/main/kafka_delay_queue_producer.go包含延迟服务的生产者:https://github.com/jiaxwu/dq/tree/main/kafka_delay_queue_example到此这篇关于Go+Kafka实现延迟消息的实现示例的文章就介绍到这了,更多相关Go Kafka延迟消息内容请搜索脚本之家以前的文章或继续浏览下面的相关文章希望大家以后多多支持脚本之家!
X 关闭
X 关闭
- 1转转集团发布2022年二季度手机行情报告:二手市场“飘香”
- 2充电宝100Wh等于多少毫安?铁路旅客禁止、限制携带和托运物品目录
- 3好消息!京东与腾讯续签三年战略合作协议 加强技术创新与供应链服务
- 4名创优品拟通过香港IPO全球发售4100万股 全球发售所得款项有什么用处?
- 5亚马逊云科技成立量子网络中心致力解决量子计算领域的挑战
- 6京东绿色建材线上平台上线 新增用户70%来自下沉市场
- 7网红淘品牌“七格格”chuu在北京又开一家店 潮人新宠chuu能红多久
- 8市场竞争加剧,有车企因经营不善出现破产、退网、退市
- 9北京市市场监管局为企业纾困减负保护经济韧性
- 10市场监管总局发布限制商品过度包装标准和第1号修改单

