如何用Golang实现Kafka生产者组去重及精确一次写入模式?
阻止Golang应用多实例向Kafka写入重复消息的实现方案
Kafka本身没有原生的「生产者组」特性来直接限制多实例的重复写入,你提到的生产者组方案本质是通过分布式协调+消息唯一性校验来实现的,下面是几种可行的落地方法,全部基于Golang实现:
1. 启用Kafka生产者幂等性(基础方案)
Kafka自0.11版本开始支持生产者幂等性,能保证同一个生产者实例对同一条消息的精确一次投递,结合全局唯一消息ID可覆盖多实例场景。
实现步骤:
- 给每条消息生成全局唯一ID(如UUID、雪花ID),作为消息的
Key或自定义Header - 配置Kafka生产者开启幂等性
Golang代码示例(使用sarama库):
package main import ( "fmt" "github.com/Shopify/sarama" "github.com/google/uuid" ) func main() { config := sarama.NewConfig() // 开启幂等性 config.Producer.Idempotent = true // 必须设置acks为all,保证消息被所有副本确认 config.Producer.RequiredAcks = sarama.WaitForAll // 多实例使用不同ClientID,避免幂等序列冲突 config.Producer.ClientID = "app-producer-" + uuid.NewString() producer, err := sarama.NewSyncProducer([]string{"kafka-broker:9092"}, config) if err != nil { panic(err) } defer producer.Close() // 生成全局唯一消息ID msgID := uuid.NewString() msg := &sarama.ProducerMessage{ Topic: "test-topic", Key: sarama.StringEncoder(msgID), Value: sarama.StringEncoder("your message content"), } partition, offset, err := producer.SendMessage(msg) if err != nil { fmt.Printf("Send message failed: %v\n", err) } else { fmt.Printf("Message sent to partition %d at offset %d\n", partition, offset) } }
2. 分布式锁控制单实例写入(针对特定业务场景)
如果你的业务要求同一业务操作只能被一个实例处理并写入Kafka,可以用分布式锁保证同一时间只有一个实例执行写入动作。
实现思路:
- 用Redis/ZooKeeper实现分布式锁,基于业务唯一标识(如订单ID、用户操作ID)作为锁的Key
- 只有获取到锁的实例才能向Kafka写入消息,写入完成后释放锁
Golang代码示例(使用Redis+redsync库):
package main import ( "fmt" "github.com/Shopify/sarama" "github.com/go-redis/redis/v8" "github.com/hashicorp/redsync/v4" "github.com/hashicorp/redsync/v4/redis/goredis/v8" "context" "time" ) func main() { // 初始化Redis客户端 client := redis.NewClient(&redis.Options{ Addr: "redis-server:6379", }) pool := goredis.NewPool(client) rs := redsync.New(pool) // 业务唯一标识,比如订单ID businessKey := "order-123456" // 创建分布式锁,设置10秒超时 mutex := rs.NewMutex(businessKey, redsync.WithExpiry(10*time.Second)) // 获取锁 if err := mutex.Lock(); err != nil { fmt.Printf("Failed to acquire lock: %v\n", err) return } defer func() { if ok, err := mutex.Unlock(); !ok || err != nil { fmt.Printf("Failed to release lock: %v\n", err) } }() // 获取锁成功,执行Kafka写入 config := sarama.NewConfig() producer, err := sarama.NewSyncProducer([]string{"kafka-broker:9092"}, config) if err != nil { panic(err) } defer producer.Close() msg := &sarama.ProducerMessage{ Topic: "order-topic", Key: sarama.StringEncoder(businessKey), Value: sarama.StringEncoder("order created"), } _, _, err = producer.SendMessage(msg) if err != nil { fmt.Printf("Send message failed: %v\n", err) } else { fmt.Println("Message sent successfully") } }
3. 基于数据库的去重校验(强一致性方案)
如果需要严格保证消息不重复写入,可以在写入Kafka前,先将消息ID写入数据库(加唯一约束),只有写入数据库成功的实例才能向Kafka发送消息。
实现思路:
- 创建去重表,包含
msg_id(唯一索引)、status等字段 - 用数据库事务执行「插入msg_id → 写入Kafka → 更新status」流程,插入失败则说明消息已存在,放弃写入
Golang代码示例(使用GORM+MySQL):
package main import ( "fmt" "github.com/Shopify/sarama" "gorm.io/driver/mysql" "gorm.io/gorm" ) // 去重表模型 type MsgDeduplicate struct { MsgID string `gorm:"primaryKey;size:64"` Status int `gorm:"default:0"` // 0:待发送 1:已发送 } func main() { // 初始化DB连接 dsn := "user:password@tcp(mysql-server:3306)/dbname?charset=utf8mb4&parseTime=True&loc=Local" db, err := gorm.Open(mysql.Open(dsn), &gorm.Config{}) if err != nil { panic(err) } // 自动建表(生产环境建议手动建表) db.AutoMigrate(&MsgDeduplicate{}) // 生成全局唯一消息ID msgID := "unique-msg-id-123" // 开启事务 tx := db.Begin() defer func() { if r := recover(); r != nil { tx.Rollback() } }() // 插入去重记录 deduplicate := MsgDeduplicate{MsgID: msgID} if err := tx.Create(&deduplicate).Error; err != nil { tx.Rollback() fmt.Printf("Message already exists: %v\n", err) return } // 写入Kafka config := sarama.NewConfig() producer, err := sarama.NewSyncProducer([]string{"kafka-broker:9092"}, config) if err != nil { tx.Rollback() panic(err) } defer producer.Close() msg := &sarama.ProducerMessage{ Topic: "test-topic", Key: sarama.StringEncoder(msgID), Value: sarama.StringEncoder("message content"), } _, _, err = producer.SendMessage(msg) if err != nil { tx.Rollback() fmt.Printf("Send message failed: %v\n", err) return } // 更新状态为已发送 if err := tx.Model(&MsgDeduplicate{}).Where("msg_id = ?", msgID).Update("status", 1).Error; err != nil { tx.Rollback() fmt.Printf("Update status failed: %v\n", err) return } // 提交事务 if err := tx.Commit().Error; err != nil { fmt.Printf("Commit transaction failed: %v\n", err) } else { fmt.Println("Message sent successfully with deduplication") } }
4. 补充:消费者端最终去重
如果写入端的去重机制有遗漏,还可以在消费者端做最终校验:
- 消费者维护本地LRU缓存或用Redis记录已消费的消息ID
- 消费消息时先检查是否已处理过,未处理再执行业务逻辑
Golang代码示例(使用LRU缓存):
package main import ( "fmt" "github.com/Shopify/sarama" "github.com/hashicorp/golang-lru/v2" ) func main() { // 初始化LRU缓存,容量10000 cache, _ := lru.New[string, bool](10000) config := sarama.NewConfig() config.Consumer.Return.Errors = true consumer, err := sarama.NewConsumer([]string{"kafka-broker:9092"}, config) if err != nil { panic(err) } defer consumer.Close() partitionConsumer, err := consumer.ConsumePartition("test-topic", 0, sarama.OffsetOldest) if err != nil { panic(err) } defer partitionConsumer.Close() for msg := range partitionConsumer.Messages() { msgID := string(msg.Key) // 检查缓存中是否存在 if _, ok := cache.Get(msgID); ok { fmt.Printf("Message %s already processed, skip\n", msgID) continue } // 处理业务逻辑 fmt.Printf("Process message: %s\n", string(msg.Value)) // 加入缓存 cache.Add(msgID, true) } }
内容的提问来源于stack exchange,提问作者Allo
相关产品推荐
相关产品推荐

