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

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)
                    );

}

已尝试方案及问题

  1. 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主题。
  2. SendToDlqAndContinue + Processor方案:Processor的process方法未被触发,原因是在process方法内嵌套了流操作,未正确处理单条记录。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 16:35:53