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
相关产品推荐
相关产品推荐

