You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 12:56:09