Spring Cloud Stream多Topic生产时DLQ、有限重试与EOS实现方案咨询
最优实现方案
推荐放弃响应式和返回Tuple的实现方式,使用命令式Consumer + StreamBridge的方案,直接复用框架内置的事务、重试、DLQ能力,逻辑简洁易维护,完全满足所有需求:
实现代码
import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.cloud.stream.function.StreamBridge; import org.springframework.transaction.annotation.Transactional; import java.util.function.Consumer; @Bean @Transactional public Consumer<Integer> numberToJson(AuditLogRepository repository, StreamBridge streamBridge) { Logger LOGGER = LoggerFactory.getLogger(getClass()); var random = new Random(); return n -> { LOGGER.info("Transforming n=" + n); var leftJson = "{ \"n\": \"" + n + "\", \"side\": \"left\" }"; var rightJson = "{ \"n\": \"" + n + "\", \"side\": \"right\" }"; // 数据库操作自动加入当前事务 repository.createIfNotExists("Transformed n=" + n); // 模拟转换异常 if (random.nextDouble() < 0.3) { LOGGER.error("Transformer failure on n=" + n); throw new RuntimeException("Transformer failure on n=" + n); } // 发送到左Topic,发送操作自动加入事务 if (!streamBridge.send("tx-json-left", leftJson)) { throw new RuntimeException("Left topic message send failed"); } // 模拟右Topic发送失败 if (random.nextDouble() < 0.1) { LOGGER.error("Failed to publish right-side JSON: n=" + n); throw new RuntimeException("Failed to publish right-side JSON: n=" + n); } // 发送到右Topic if (!streamBridge.send("tx-json-right", rightJson)) { throw new RuntimeException("Right topic message send failed"); } }; }
配置文件
spring: cloud: stream: kafka: binder: transaction: transaction-id-prefix: 'tx-' producer: configuration: retries: 3 acks: all bindings: numberToJson-in-0: consumer: enableDlq: true dlqName: error.tx-number.numberToJson bindings: numberToJson-in-0: destination: tx-number group: numberToJson consumer: maxAttempts: 3 # 全流程最多重试3次 properties: isolation.level: read_committed function: definition: numberToJson
方案符合要求验证
- 单输入转双输出:从
tx-number单Topic消费消息,通过StreamBridge分别投递到左右两个目标Topic - 原子性保证:
@Transactional注解将消费位点提交、两次消息投递、数据库操作包裹在同一个Kafka事务中,任意环节抛出异常会触发全事务回滚,不会出现单边投递成功的问题 - 重试控制:框架内置的
maxAttempts=3配置实现全流程最多3次重试,不需要手动编写重试逻辑 - DLQ自动投递:开启
enableDlq后,3次重试全部失败的消息会被框架自动投递到DLQ Topic,不会阻塞后续消息消费
原有方案问题说明
- 命令式返回Tuple2启动报错:你使用了reactor的
reactor.util.function.Tuple2,Spring Cloud Stream要求使用org.springframework.cloud.function.context.Tuple才能支持多输出绑定,切换包后还需要额外配置输出绑定映射,事务集成的灵活性远低于StreamBridge方案 - 响应式方案存在的缺陷:自行维护Sink和重试逻辑会导致事务边界模糊,响应式场景下消费位点提交默认和事务生命周期不同步,手动写的多层retry逻辑会导致事务重复开启,极易出现消息重复投递、丢失的问题
内容的提问来源于stack exchange,提问作者Tomboyo
相关产品推荐
相关产品推荐

