如何处理Kafka生产者发送与数据库存储的事务一致性问题
实现方案
核心逻辑
利用Kafka生产者的消息确认回调机制,只有收到Kafka集群返回的发送成功响应后,再执行PostgreSQL的写入操作,从流程上保证两个操作的先后顺序。
具体实现步骤
- 先配置Kafka生产者的可靠性参数,避免消息发送成功的误判:
- 设置
acks为all(或-1),要求消息被ISR列表中所有副本同步后才返回成功响应 - 配置合理的
retries重试次数,规避临时网络波动导致的发送失败 - 可选开启
idempotence幂等配置,避免重试带来的重复消息问题
- 设置
- 序列化待发送的业务对象为Kafka支持的格式(JSON、Avro、Protobuf均可),构造ProducerRecord对象
- 采用带回调的异步发送方式提交消息到Kafka,在回调逻辑中判断发送结果:
- 无异常返回则判定为发送成功,执行PostgreSQL的对象写入逻辑
- 有异常则判定为发送失败,记录日志后执行降级逻辑,比如存入本地重试表后续补偿
注意:如果需要严格保证Kafka发送和PG写入的最终一致性,建议新增本地状态表:发送Kafka前先写入状态为「待处理」的记录,Kafka发送成功后更新为「已发待存库」,PG写入成功后更新为「已完成」,新增定时任务扫描超时的异常记录做补偿处理。
代码示例(Java版本参考)
Kafka生产者配置
Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-node1:9092,kafka-node2:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class.getName()); // 可靠性配置 props.put(ProducerConfig.ACKS_CONFIG, "all"); props.put(ProducerConfig.RETRIES_CONFIG, 3); props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); KafkaProducer<String, BusinessObject> producer = new KafkaProducer<>(props);
核心业务逻辑
// 待处理的业务对象 BusinessObject data = buildBusinessData(); ProducerRecord<String, BusinessObject> record = new ProducerRecord<>("your-topic-name", data.getBizId(), data); // 异步发送带回调 producer.send(record, (metadata, exception) -> { if (exception == null) { // Kafka发送成功,写入PostgreSQL,建议加重试逻辑 try { saveToPg(data); } catch (SQLException e) { // 写入PG失败,记录日志+告警,后续人工补偿或者异步重试 log.error("PG写入失败,业务ID:{}", data.getBizId(), e); } } else { // Kafka发送失败,执行降级逻辑 log.error("Kafka发送失败,业务ID:{}", data.getBizId(), exception); } });
低并发场景同步实现
如果业务并发量低,也可以用同步发送的方式简化逻辑:
try { // 阻塞等待Kafka返回发送结果 producer.send(record).get(3, TimeUnit.SECONDS); // 发送成功后写入PG saveToPg(data); } catch (Exception e) { log.error("操作失败,业务ID:{}", data.getBizId(), e); }
内容的提问来源于stack exchange,提问作者foreigninstuction
相关产品推荐
相关产品推荐

