Spring Cloud中HTTP消息转Kafka消息发送失败问题排查
问题:Spring Cloud Stream + Kafka + HTTP响应式消息发送失败排查
我在Spring环境中尝试结合spring-cloud-stream-binder-kafka(Kafka流处理)与spring-cloud-function-web(HTTP调用),以响应式方式将HTTP接收的消息发送到Kafka主题。目前HTTP层能成功接收消息,但消息并未发送到配置的Kafka主题,需要排查原因。
1. 核心代码实现
Event消费类(接收HTTP消息)
import org.springframework.messaging.Message import org.springframework.messaging.support.MessageBuilder import org.springframework.transaction.annotation.Transactional import reactor.core.publisher.Flux import reactor.core.publisher.Sinks class Event : Consumer<Flux<Message<String>>> { protected val logger = loggerFor(javaClass) protected val unicastProcessor = Sinks.many().unicast().onBackpressureBuffer<Message<String>>() @Transactional override fun accept(t: Flux<Message<String>>) { t.map { msg -> val txt = msg.payload.let { MessageBuilder.withPayload(it).build() } logger.info("Received HTTP Message: $txt") unicastProcessor.emitNext(msg, Sinks.EmitFailureHandler.FAIL_FAST) } .subscribe() } }
生产者Bean配置
@Bean fun produce(): Supplier<Flux<Message<String>>> = Supplier { unicastProcessor.asFlux() }
2. 配置文件(application.yml)
cloud: function: scan: packages: com.example.blabla.functions stream: bindings: produce-out-0: destination: stats-topic binder: kafka consume-in-0: destination: stats-topic binder: kafka group: stats default-binder: kafka kafka: binder: auto-create-topics: true auto-create-partitions: true brokers: "kafka:9092" consumer-properties: key.deserializer: org.apache.kafka.common.serialization.StringDeserializer value.deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer producer-properties: key.serializer: org.apache.kafka.common.serialization.StringSerializer value.serializer: io.confluent.kafka.serializers.KafkaAvroSerializer transaction: transactionIdPrefix: transactionPrefixTest
3. 当前现象
HTTP层已成功接收消息,日志如下:
2024-09-17T16:02:39.477Z INFO 6010 --- [jvmtask] [or-http-epoll-2] com.example.blabla.functions.Event : Received HTTP Message: GenericMessage [payload=text to log 3, headers={id=6b339ecf-d62f-fe28-f7b9-10520587a66e, timestamp=1726588959474}]
但消息未发送到Kafka的stats-topic主题。
排查方向
- Sink类型不匹配:当前使用
unicast()类型的Sink,该类型仅支持单个订阅者。若Spring Cloud Stream的Supplier订阅时机滞后,或存在其他订阅者,会导致消息无法传递到生产者。建议替换为multicast()类型:Sinks.many().multicast().onBackpressureBuffer()。 - 事务注解误用:
accept方法添加了@Transactional,但响应式流的事务需配合Reactive事务管理器(如@Transactional(transactionManager = "reactiveTransactionManager"))。配置不当会阻塞流处理或导致消息无法提交,可先移除事务注解测试。 - 序列化不兼容:生产者配置使用
KafkaAvroSerializer,但发送的是Message<String>,String类型无法直接用Avro序列化。要么将消息转为Avro对象,要么将生产者的value.serializer改为StringSerializer。 - 实例引用不一致:
unicastProcessor是Event类的protected成员,若Event未声明为Spring Bean(如未加@Component),配置类中的Supplier可能引用了不同的Event实例,导致Sink无法共享。需确保Event被Spring管理,并在配置类中注入Event实例获取Sink。 - 绑定配置错误:检查
produce-out-0的绑定是否与Supplier函数名(produce)对应,同时确认function.scan.packages包含Event类和Supplier Bean的包路径。 - 日志调试:开启Kafka Binder的DEBUG日志,查看生产者初始化、消息发送的详细日志,排查是否存在序列化错误、事务提交失败等异常:
logging: level: org.springframework.cloud.stream.binder.kafka: DEBUG org.apache.kafka: DEBUG
内容的提问来源于stack exchange,提问作者Mostafa Abdelhamid
相关产品推荐
相关产品推荐

