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

Spring Cloud Streams结合Kafka Streams Binder函数式处理器异常处理示例求问

Spring Cloud Stream Kafka Streams Binder 函数式风格处理器运行时异常DLQ完整实现

1. 依赖引入

Maven pom.xml 核心依赖配置:

<dependencies>
    <dependency>
        <groupId>org.springframework.cloud</groupId>
        <artifactId>spring-cloud-stream-binder-kafka-streams</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-actuator</artifactId>
        <optional>true</optional>
    </dependency>
</dependencies>

依赖版本需和当前使用的Spring Boot、Spring Cloud版本基线匹配。

2. 应用配置

application.yml 配置示例(适配多处理器场景):

spring:
  cloud:
    function:
      # 多处理器用分号分隔定义
      definition: orderProcess; payProcess
    stream:
      bindings:
        # 第一个处理器orderProcess的绑定
        orderProcess-in-0:
          destination: order-topic
          group: order-consumer-group
        orderProcess-out-0:
          destination: order-after-process-topic
        # 第二个处理器payProcess的绑定
        payProcess-in-0:
          destination: pay-topic
          group: pay-consumer-group
        payProcess-out-0:
          destination: pay-after-process-topic
        # 通用DLQ输出绑定
        dlq-out-0:
          destination: common-dlq-topic
      kafka:
        streams:
          binder:
            brokers: your-kafka-broker:9092
            configuration:
              # 已有反序列化异常DLQ配置保留
              default.deserialization.exception.handler: org.springframework.kafka.streams.RecoveringDeserializationExceptionHandler
              spring.json.trusted.packages: "*"
            dlqProducerProperties:
              key.serializer: org.apache.kafka.common.serialization.StringSerializer
              value.serializer: org.apache.kafka.common.serialization.StringSerializer

3. 通用异常处理工具封装

封装统一的业务异常捕获逻辑,避免每个处理器重复写try-catch:

import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.kstream.KStream;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.stream.function.StreamBridge;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

@Component
public class DlqProcessWrapper {
    @Autowired
    private StreamBridge streamBridge;

    public <K, V> KStream<K, V> wrapProcess(KStream<K, V> input, String processorName, java.util.function.BiFunction<K, V, V> processLogic) {
        return input.map((key, value) -> {
            try {
                V result = processLogic.apply(key, value);
                return KeyValue.pair(key, result);
            } catch (Exception e) {
                // 异常消息发送到DLQ
                streamBridge.send("dlq-out-0", MessageBuilder.withPayload(value)
                        .setHeader("dlq-origin-topic", processorName)
                        .setHeader("dlq-exception-type", e.getClass().getName())
                        .setHeader("dlq-exception-message", e.getMessage())
                        .build());
                // 异常消息跳过后续业务处理,返回null会被自动过滤
                return null;
            }
        }).filter((key, value) -> value != null);
    }
}

4. 多处理器业务实现

基于封装工具实现多处理器的业务逻辑,无需重复处理异常:

import org.apache.kafka.streams.kstream.KStream;
import org.springframework.context.annotation.Bean;
import org.springframework.stereotype.Component;
import java.util.function.Function;

@Component
public class MultiProcessorConfig {
    private final DlqProcessWrapper dlqProcessWrapper;

    public MultiProcessorConfig(DlqProcessWrapper dlqProcessWrapper) {
        this.dlqProcessWrapper = dlqProcessWrapper;
    }

    @Bean
    public Function<KStream<String, Order>, KStream<String, Order>> orderProcess() {
        return input -> dlqProcessWrapper.wrapProcess(input, "order-topic", (key, order) -> {
            // 实际业务处理逻辑,抛出异常会自动进入DLQ
            if (order.getAmount() == null) {
                throw new IllegalArgumentException("订单金额不能为空");
            }
            order.setStatus("PROCESSED");
            return order;
        });
    }

    @Bean
    public Function<KStream<String, PayRecord>, KStream<String, PayRecord>> payProcess() {
        return input -> dlqProcessWrapper.wrapProcess(input, "pay-topic", (key, payRecord) -> {
            // 实际支付处理逻辑
            if (payRecord.getPayId() == null) {
                throw new IllegalArgumentException("支付ID不能为空");
            }
            payRecord.setPaySuccess(true);
            return payRecord;
        });
    }
}

5. 效果验证

  • 向order-topic发送金额为空的订单消息,异常消息会自动进入common-dlq-topic,同时携带异常来源、异常类型等头信息
  • 正常消息会进入后续的order-after-process-topic进行下一步流转
  • 可通过/actuator/kafkastreams/topology端点查看最终生成的流处理拓扑,确认处理链路符合预期

内容的提问来源于stack exchange,提问作者gira1

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 12:36:00