尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Go消息队列:Kafka与NATS集成

Go消息队列:Kafka与NATS集成 Go消息队列:Kafka与NATS集成摘要: 本篇讲解Go语言集成Kafka和NATS两种消息队列用sarama库实现Kafka生产者和消费者组用nats.go实现NATS发布订阅分享消息幂等处理的去重方案分析Kafka自动提交offset导致消息丢失的踩坑经验对比Kafka、NATS和RabbitMQ三种消息队列。开篇故事我们的订单事件流原先是同步写数据库的订单创建、支付、发货各写一张表。业务量上来后数据库扛不住高峰期订单创建的写入QPS到2000主库CPU飙到90%。架构组决定上消息队列做异步解耦订单服务发事件下游库存、物流、通知各消费各的。选型时纠结过Kafka和NATS。Kafka吞吐高、有持久化、生态成熟但运维重要Zookeeper(新版KRaft好一些)。NATS轻量、延迟低、部署简单单机就能跑。最后核心订单流走Kafka保可靠性内部微服务通信走NATS求轻快。上线第一周踩了个坑。消费者用自动提交offset服务重启时有200多条消息没处理完就提交了直接丢消息。后来改成手动提交加幂等消费才解决。这篇把Kafka和NATS的集成写清楚。一、Kafka生产者与消费者Kafka的Go客户端主流用sarama库IBM维护功能完整。生产者发消息到topic消费者组从topic拉取消息每个分区同一时刻只有一个消费者消费保证分区有序。先看生产者。sarama的同步生产者发完等确认可靠性高但吞吐低。异步生产者发完不等确认吞吐高但有丢消息风险。订单事件不能丢用同步生产者。packagemqimport(contextencoding/jsonlogtimegithub.com/IBM/sarama)// KafkaProducer Kafka同步生产者typeKafkaProducerstruct{producer sarama.SyncProducer// 同步生产者topicstring// 默认topic}// NewKafkaProducer 创建Kafka生产者// brokers: Kafka broker地址列表// topic: 默认发送的topicfuncNewKafkaProducer(brokers[]string,topicstring)(*KafkaProducer,error){// 配置生产者config:sarama.NewConfig()// 等待所有副本确认最高可靠性config.Producer.RequiredAckssarama.WaitForAll// 启用成功发送的确认通知config.Producer.Return.Successestrue// 启用错误返回config.Producer.Return.Errorstrue// 消息压缩节省带宽config.Producer.Compressionsarama.CompressionSnappy// 幂等生产者防止重试导致重复config.Producer.Idempotenttrue// 幂等模式要求等所有副本确认config.Net.MaxOpenRequests1// 创建同步生产者producer,err:sarama.NewSyncProducer(brokers,config)iferr!nil{returnnil,err}returnKafkaProducer{producer:producer,topic:topic,},nil}// Send 发送消息// key: 分区key相同key的消息进同一分区// value: 消息内容会被JSON序列化func(p*KafkaProducer)Send(ctx context.Context,keystring,valueinterface{})error{// 序列化消息体data,err:json.Marshal(value)iferr!nil{returnerr}// 构造Kafka消息msg:sarama.ProducerMessage{Topic:p.topic,Key:sarama.StringEncoder(key),Value:sarama.ByteEncoder(data),Timestamp:time.Now(),}// 同步发送等待broker确认partition,offset,err:p.producer.SendMessage(msg)iferr!nil{returnerr}log.Printf(消息发送成功: topic%s partition%d offset%d,p.topic,partition,offset)returnnil}// Close 关闭生产者刷新缓冲区func(p*KafkaProducer)Close()error{returnp.producer.Close()}再看消费者组。消费者组是多消费者协作消费一个topicKafka自动分配分区每个分区内消息有序。关键配置是offset提交策略。packagemqimport(contextloggithub.com/IBM/sarama)// OrderHandler 订单消息处理器// 实现sarama.ConsumerGroupHandler接口typeOrderHandlerstruct{processorfunc(msg[]byte)error// 业务处理函数}// NewOrderHandler 创建消息处理器// processor: 实际处理消息的回调函数funcNewOrderHandler(processorfunc([]byte)error)*OrderHandler{returnOrderHandler{processor:processor}}// Setup 消费者组会话开始时调用func(h*OrderHandler)Setup(sarama.ConsumerGroupSession)error{// 可以在这里做初始化比如加载缓存returnnil}// Cleanup 消费者组会话结束时调用func(h*OrderHandler)Cleanup(sarama.ConsumerGroupSession)error{// 可以在这里做清理returnnil}// ConsumeClaim 核心消费逻辑每条消息调用一次func(h*OrderHandler)ConsumeClaim(session sarama.ConsumerGroupSession,claim sarama.ConsumerGroupClaim,)error{// 遍历channel里的消息formsg:rangeclaim.Messages(){// 调用业务处理函数err:h.processor(msg.Value)iferr!nil{// 处理失败记录日志不提交offset// 下次启动会重新消费这条消息log.Printf(消息处理失败: offset%d err%v,msg.Offset,err)returnerr}// 处理成功手动提交offsetsession.MarkMessage(msg,)}returnnil}// KafkaConsumer Kafka消费者组封装typeKafkaConsumerstruct{group sarama.ConsumerGroup topicstring}// NewKafkaConsumer 创建消费者组// brokers: broker地址// groupID: 消费者组名// topic: 消费的topicfuncNewKafkaConsumer(brokers[]string,groupID,topicstring)(*KafkaConsumer,error){config:sarama.NewConfig()// 关闭自动提交改用手动提交// 自动提交可能在消息未处理完就提交重启丢消息config.Consumer.Offsets.AutoCommit.Disabletrue// 从最早的offset开始消费config.Consumer.Offsets.Initialsarama.OffsetOldest// 消费者组需要启用Return.Errorsconfig.Consumer.Return.Errorstruegroup,err:sarama.NewConsumerGroup(brokers,groupID,config)iferr!nil{returnnil,err}returnKafkaConsumer{group:group,topic:topic},nil}// Consume 开始消费阻塞直到context取消// handler: 消息处理器func(c*KafkaConsumer)Consume(ctx context.Context,handler sarama.ConsumerGroupHandler)error{for{// ConsumeClaim会阻塞发生rebalance后返回err:c.group.Consume(ctx,[]string{c.topic},handler)iferr!nil{returnerr}// 检查context是否已取消ifctx.Err()!nil{returnctx.Err()}// rebalance后重新进入消费循环}}// Close 关闭消费者func(c*KafkaConsumer)Close()error{returnc.group.Close()}手动提交offset的核心在session.MarkMessage。这条消息处理成功才标记消费者组会话提交时只提交已标记的消息。处理失败的返回错误不标记重启后重新消费。二、NATS发布订阅NATS的Go客户端是nats.goAPI比sarama简洁。NATS核心模式是pub/sub发消息不持久化订阅者不在线就丢消息。需要持久化用NATS JetStream带消息存储和确认机制。内部服务通信用核心pub/sub够用。比如用户服务发布用户创建事件通知服务和积分服务各自订阅。packagemqimport(contextencoding/jsonlogtimegithub.com/nats-io/nats.go)// NatsConn NATS连接封装typeNatsConnstruct{nc*nats.Conn// NATS原生连接}// NewNatsConn 创建NATS连接// url: NATS服务器地址如nats://127.0.0.1:4222funcNewNatsConn(urlstring)(*NatsConn,error){// 连接选项opts:[]nats.Option{// 连接超时nats.Timeout(5*time.Second),// 重连nats.ReconnectWait(2*time.Second),// 最大重连次数-1表示无限重连nats.MaxReconnects(-1),// 重连成功回调nats.ReconnectHandler(func(nc*nats.Conn){log.Printf(NATS重连成功: %s,nc.ConnectedUrl())}),// 断线回调nats.DisconnectErrHandler(func(nc*nats.Conn,errerror){log.Printf(NATS断开: %v,err)}),}nc,err:nats.Connect(url,opts...)iferr!nil{returnnil,err}returnNatsConn{nc:nc},nil}// Publish 发布消息// subject: 主题类似Kafka的topic// data: 消息内容自动JSON序列化func(c*NatsConn)Publish(subjectstring,datainterface{})error{// 序列化消息payload,err:json.Marshal(data)iferr!nil{returnerr}// 发布到指定subjectreturnc.nc.Publish(subject,payload)}// Subscribe 订阅消息// subject: 订阅的主题// handler: 消息处理回调// 返回取消订阅函数func(c*NatsConn)Subscribe(subjectstring,handlerfunc([]byte)error)(*nats.Subscription,error){returnc.nc.Subscribe(subject,func(msg*nats.Msg){// 调用业务处理函数iferr:handler(msg.Data);err!nil{log.Printf(NATS消息处理失败: subject%s err%v,subject,err)}})}// Request 请求-响应模式// 发送请求并等待响应适合RPC场景func(c*NatsConn)Request(ctx context.Context,subjectstring,data[]byte,timeout time.Duration)([]byte,error){// 带context的请求支持超时取消resp,err:c.nc.RequestWithContext(ctx,subject,data)iferr!nil{returnnil,err}returnresp.Data,nil}// Close 关闭连接func(c*NatsConn)Close(){c.nc.Close()}NATS还有队列订阅(queue group)模式同一队列组的多个订阅者竞争消费每条消息只投递给组内一个订阅者。这个模式和Kafka消费者组类似适合多实例分摊消费。// QueueSubscribe 队列订阅// 同一queueGroup的订阅者竞争消费func(c*NatsConn)QueueSubscribe(subject,queueGroupstring,handlerfunc([]byte)error)(*nats.Subscription,error){returnc.nc.QueueSubscribe(subject,queueGroup,func(msg*nats.Msg){iferr:handler(msg.Data);err!nil{log.Printf(队列消费失败: subject%s err%v,subject,err)}})}三、消息幂等处理消息队列的消费者可能收到重复消息。网络重试、消费者重启、rebalance都会导致重复投递。业务层必须做幂等处理处理多次和一次效果一样。幂等最常用的方案是基于业务唯一键去重。用Redis记录已处理的消息ID消费前先检查。packagemqimport(contextgithub.com/redis/go-redis/v9)// IdempotentProcessor 幂等消息处理器typeIdempotentProcessorstruct{rdb*redis.Client// Redis客户端存已处理消息IDttlint// 去重记录保留时长单位秒}// NewIdempotentProcessor 创建幂等处理器// rdb: Redis客户端// ttl: 去重记录保留秒数建议大于消息最大重试间隔funcNewIdempotentProcessor(rdb*redis.Client,ttlint)*IdempotentProcessor{returnIdempotentProcessor{rdb:rdb,ttl:ttl,}}// Process 处理消息自动去重// msgKey: 业务唯一键如订单号// handler: 实际业务处理回调func(p*IdempotentProcessor)Process(ctx context.Context,msgKeystring,handlerfunc()error)error{// Lua脚本: SETNX检查成功则标记luaScript: if redis.call(SETNX, KEYS[1], ARGV[1]) 1 then redis.call(EXPIRE, KEYS[1], ARGV[2]) return 1 else return 0 end result,err:p.rdb.Eval(ctx,luaScript,[]string{msg:msgKey},1,p.ttl).Int()iferr!nil{returnerr}ifresult0{// 已处理过幂等跳过returnnil}// 执行业务逻辑iferr:handler();err!nil{// 失败则删除标记允许下次重试p.rdb.Del(ctx,msg:msgKey)returnerr}returnnil}关键点是处理失败要删掉占位key否则这条消息永远不会再处理变成死信。去重记录的TTL要大于消息最大重试间隔否则TTL过期后重复消息又会处理一次。四、踩坑经验:Kafka自动提交offset导致消息丢失这个坑前面提过详细讲一下。消费者组配置了config.Consumer.Offsets.AutoCommit.Enable true默认每秒自动提交一次offset。我们的订单消息处理涉及调外部支付接口单条耗时800毫秒到2秒不等。某次发布重启服务有200多条消息已经从Kafka拉到内存但没处理完。自动提交把offset推到了这批消息的末尾。服务重启后从新offset消费这200多条消息全丢了。问题根源是自动提交不等消息处理完就提交offset。offset只表示拉取到哪了不表示处理完到哪了。解决方案有两个层面。第一关闭自动提交改手动提交前面代码里已经做了。第二加幂等处理兜底万一offset提交错误重复消费不会造成数据问题。手动提交的代码前面已经写了session.MarkMessage只在处理成功后调用。但手动提交有个细节要注意MarkMessage只是标记真正的提交发生在会话结束或定时触发时。如果服务突然崩溃标记了但没提交的消息还是会重复消费。所以幂等处理是最终兜底。// 消费者主循环演示手动提交的使用方式funcconsumeOrders(ctx context.Context,consumer*KafkaConsumer,rdb*redis.Client){// 创建幂等处理器TTL设为1小时processor:NewIdempotentProcessor(rdb,3600)// 创建消息处理器handler:NewOrderHandler(func(msg[]byte)error{// 用订单号作为幂等键// 实际从消息体解析订单号msgKey:order_12345returnprocessor.Process(ctx,msgKey,func()error{// 执行实际业务逻辑: 反序列化、更新库存returnnil})})// 开始消费_consumer.Consume(ctx,handler)}两层保障手动提交减少重复消费概率幂等处理保证重复消费不产生副作用。线上跑了一年多偶尔有重复消费没再出现丢消息。五、对比分析特性KafkaNATS CoreNATS JetStream持久化有无有吞吐极高(10万/s)高(百万/s)中延迟中(毫秒级)极低(亚毫秒)低消费者组内置支持队列组内置支持运维复杂度高低中消息确认offset提交无ACK确认适合场景事件流、日志内部通信、RPC可靠投递Kafka吞吐和持久化强适合核心业务事件流和日志收集。NATS Core延迟极低适合内部服务通信和请求-响应模式。NATS JetStream介于两者之间需要持久化但不想运维Kafka的场景可以用。选型看场景没有一劳永逸的方案。总结Kafka用sarama集成生产者开幂等和全副本确认保证可靠性消费者组关闭自动提交改手动提交。NATS用nats.go集成核心pub/sub轻量低延迟JetStream补上持久化能力。消息消费必须做幂等处理Redis去重加Lua原子操作是常用方案。手动提交offset配幂等处理是消息不丢不重的双保险。
返回列表