Spring Cloud Stream Kafka Streams函数无需Processor实现消息转DLQ遇阻
Spring Cloud Stream Kafka Streams 异常消息转DLQ实现
需求
在Spring Cloud Stream Kafka Streams的Function中处理消息时,将setSomething方法执行时因逻辑规则抛出异常的消息发送至DLQ。
现有配置与代码
配置文件(application.yaml)
function: definition: filterConsumption|EnrichConsumption; #Functions and topic binding stream: bindings: filterConsumptionEnrichConsumption-in-0: destination: input filterConsumptionEnrichConsumption-out-0: destination: output kafka: streams: bindings: filterConsumptionEnrichConsumption-in-0: consumer: enable-dlq: true dlqName: input_dlq application-id: input-application-id filterConsumptionEnrichConsumption-out-1: consumer: enable-dlq: false application-id: output-application-id binder: #Kafka consumer config replicationFactor: ${KAFKA_STREAM_REPLICATION_FACTOR:1} brokers: ${BOOTSTRAP_SERVERS_CONFIG:localhost:9092} deserialization-exception-handler: sendToDlq
函数代码
@Bean("EnrichConsumption") public Function<KStream<String, ConsumptionSchema>, KStream<String, ConsumptionSchema>> EnrichConsumption() { return input -> input.filter((key, consumptions) -> !getSomething(consumptions).orElse("").isBlank()) .merge( //filter consumptions having a tradingName input.filter((key, consumptions) -> getSomething(consumptions).orElse("").isBlank()) //enrich consumptions with missing tradingName .mapValues(this::setSomething) ); }
已尝试方案及问题
- StreamBridge方案:抛出Kafka连接错误
Connection to node -1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available.,即使已通过BOOTSTRAP_SERVERS_CONFIG配置实际Kafka地址,正常消息可路由至output主题。 - SendToDlqAndContinue + Processor方案:Processor的
process方法未被触发,原因是在process方法内嵌套了流操作,未正确处理单条记录。 - DltAwareProcessor方案:处理器被正常调用,但执行
streamBridge.send时出现相同的Kafka连接错误。
解决方案
方案一:使用SendToDlqAndContinue(优先推荐)
问题出在之前的Processor实现错误,不应在process方法内重复操作原始流,而是要针对每条记录单独处理,同时正确利用Kafka Streams的ProcessorContext和SendToDlqAndContinue工具类。
修正后的代码:
@Autowired private SendToDlqAndContinue dlqHandler; private static final Logger log = LoggerFactory.getLogger(YourClassName.class); @Bean("EnrichConsumption") public Function<KStream<String, ConsumptionSchema>, KStream<String, ConsumptionSchema>> EnrichConsumption() { return input -> { // 处理已有有效字段的记录 KStream<String, ConsumptionSchema> validRecords = input .filter((key, consumptions) -> !getSomething(consumptions).orElse("").isBlank()); // 处理需要补全字段的记录,异常时发送DLQ KStream<String, ConsumptionSchema> enrichedRecords = input .filter((key, consumptions) -> getSomething(consumptions).orElse("").isBlank()) .process(() -> new Processor<String, ConsumptionSchema>() { private ProcessorContext context; @Override public void init(ProcessorContext context) { this.context = context; } @Override public void process(String key, ConsumptionSchema value) { try { ConsumptionSchema enriched = setSomething(value); context.forward(key, enriched); } catch (DlqException e) { log.error("处理消息异常,key: {}", key, e); // 构造符合要求的ConsumerRecord用于DLQ发送 ConsumerRecord<String, ConsumptionSchema> consumerRecord = new ConsumerRecord<>( context.topic(), context.partition(), context.offset(), key, value ); dlqHandler.sendToDlq(consumerRecord, e); // 异常消息不转发,避免进入下游流 } } @Override public void close() { // 按需清理资源 } }); return validRecords.merge(enrichedRecords); }; }
配置说明:
确保input_dlq主题已提前创建,或开启Kafka自动创建主题配置(auto.create.topics.enable=true)。同时确认spring.cloud.stream.kafka.streams.bindings.filterConsumptionEnrichConsumption-in-0.consumer.enable-dlq=true配置生效,该配置会自动为输入通道绑定DLQ处理器。
方案二:修复StreamBridge连接问题
StreamBridge出现连接错误,是因为其使用的Kafka生产者配置未正确读取BOOTSTRAP_SERVERS_CONFIG,需要显式配置StreamBridge对应的生产者绑定:
1. 补充配置
在application.yaml中添加DLQ输出绑定配置:
stream: bindings: # 新增DLQ输出绑定 dlq-out-0: destination: input_dlq kafka: bindings: dlq-out-0: producer: configuration: bootstrap.servers: ${BOOTSTRAP_SERVERS_CONFIG:localhost:9092}
2. 修正后的StreamBridge代码
private final StreamBridge streamBridge; private static final Logger log = LoggerFactory.getLogger(YourClassName.class); public YourClassName(StreamBridge streamBridge) { this.streamBridge = streamBridge; } @Bean("EnrichConsumption") public Function<KStream<String, ConsumptionSchema>, KStream<String, ConsumptionSchema>> EnrichConsumption() { return input -> { KStream<String, ConsumptionSchema> validRecords = input .filter((key, consumptions) -> !getSomething(consumptions).orElse("").isBlank()); KStream<String, ConsumptionSchema> enrichedRecords = input .filter((key, consumptions) -> getSomething(consumptions).orElse("").isBlank()) .mapValues((key, value) -> { try { return setSomething(value); } catch (DlqException e) { log.error("处理消息异常,key: {}", key, e); // 使用配置好的dlq-out-0通道发送DLQ streamBridge.send("dlq-out-0", value); return null; } }) .filter((key, value) -> value != null); // 过滤空值,避免下游处理无效消息 return validRecords.merge(enrichedRecords); }; }
说明:
- 为StreamBridge创建单独的输出绑定
dlq-out-0,显式指定Kafka broker地址,确保生产者配置正确加载。 - 发送DLQ后返回null,通过
filter移除空值,避免无效消息进入下游流。
内容的提问来源于stack exchange,提问作者Romain Huber
相关产品推荐
相关产品推荐

