如何有序消费Kafka消息避免同uuid状态异常及现有方案优化
根因分析
- Kafka生产未绑定分区:如果发送消息时未将
uuid作为消息key,同uuid的消息会被哈希分配到不同分区,多消费者并发消费不同分区时自然会出现乱序 - 消费侧处理逻辑无序:即使Kafka分区内消息有序,现有代码将所有消息写入全局
msgQueue后,如果是多协程从队列取数据消费,也会出现后到的消息先完成MySQL写入的情况 - 现有方案逻辑缺陷:仅当消息
status != 1时才添加冲突更新逻辑,status = 1的消息走普通INSERT逻辑,当对应uuid已存在(已提前写入status=3)时就会触发唯一键冲突报错。
优化方案
方案1:最小改动修复写入逻辑(推荐,无需修改上下游链路)
直接删除原有的if data.Status != 1判断,统一加上冲突更新逻辑,同时添加状态校验,仅当新状态大于已存储状态时才更新,既不会出现唯一键报错,也不会出现老状态覆盖新状态的问题,代码示例如下:
conn := &gorm.DB{} data := &Log{} // 所有写入统一走冲突更新逻辑,仅当新状态更大时才更新 conn = conn.Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "uuid"}}, DoUpdates: clause.Assignments(map[string]interface{}{ "status": gorm.Expr("IF(VALUES(status) > status, VALUES(status), status)"), // 如果有其他字段需要同步更新也可以加在这里,同样可以加条件判断 }), }) if err := conn.Create(data).Error; err != nil { return err } return nil
方案2:从根源保证消息顺序(适合对性能要求高,希望减少无效写入的场景)
- 生产侧改造:发送Kafka消息时指定消息key为
uuid,Kafka会按key哈希将同uuid的所有消息分配到同一个分区,保证分区内消息严格按生产顺序排列 - 消费侧改造:取消全局
msgQueue多协程消费的模式,改为单协程按分区顺序消费处理,或者按uuid哈希将消息分配到固定的协程处理,保证同uuid的消息永远按顺序被处理,从根源避免乱序问题。
可选补充
如果消息携带生产时间戳,也可以将状态判断替换为时间戳判断,保证永远用最新生产的消息覆盖老消息,逻辑和上述状态判断一致。
内容的提问来源于stack exchange,提问作者CharmCcc
相关产品推荐
相关产品推荐

