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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 09:35:57