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

Spring Cloud Stream集成Kafka按MessageKey分区消费异常排查

问题根因

所有消息全部路由到分区0的核心原因有3个:

  • Spring Cloud Stream Kafka Binder默认不会自动使用KafkaHeaders.MESSAGE_KEY做分区路由,未显式配置分区规则时,会默认将所有消息固定发送到分区0,和你是否设置消息key无关。
  • 生产者配置中partition-count设为3,和你实际规划的2分区主题不一致,分区数不匹配会进一步干扰路由逻辑。
  • 消费端Json反序列化未配置信任包,后续会触发类型转换错误,属于隐藏配置问题。
修正方案

1. 调整生产者配置

核心是显式开启按消息key分区的规则,同时对齐分区数配置,修改后的生产者application.yml配置如下:

spring:
  cloud:
    stream:
      kafka:
        binder:
          replicationFactor: 2
          auto-create-topics: true
          brokers: localhost:9092,localhost:9093,localhost:9094
          auto-add-partitions: true
        bindings:
          simulatePf-out-0:
            producer:
              configuration:
                key.serializer: org.apache.kafka.common.serialization.StringSerializer
                value.serializer: org.springframework.kafka.support.serializer.JsonSerializer
                # 可选:使用Kafka原生默认分区器,按key哈希路由
                partitioner.class: org.apache.kafka.clients.producer.internals.DefaultPartitioner
              # 核心配置:指定用Kafka消息key作为分区计算依据
              partition-key-expression: headers['kafka_messageKey']
      bindings:
        simulatePf-out-0:
          producer:
            useNativeEncoding: true
            # 分区数和主题实际规划的2个分区对齐
            partition-count: 2
          destination: pf-topic
          content-type: text/plain
          group: dsa-back-end

说明:配置partition-key-expression后,Spring Cloud Stream会提取消息头中的kafka_messageKey(即你代码里设置的KafkaHeaders.MESSAGE_KEY)作为分区键,相同key的消息会固定路由到同一个分区。如果配置了原生DefaultPartitioner,分区计算逻辑完全和原生Kafka客户端一致。

2. 调整消费端配置

补全Json反序列化信任配置,显式对齐分区数,修改后的消费端application.yml配置如下:

spring:
  cloud:
    stream:
      kafka:
        binder:
          replicationFactor: 2
          auto-create-topics: true
          brokers: localhost:9092,localhost:9093,localhost:9094
          min-partition-count: 2
        bindings:
          simulatePf-in-0:
            consumer:
              configuration:
                key.deserializer: org.apache.kafka.common.serialization.StringDeserializer
                value.deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
                # 配置Json反序列化信任所有包,避免类型转换报错
                spring.json.trusted.packages: "*"
      bindings:
        simulatePf-in-0:
          destination: pf-topic
          content-type: text/plain
          group: powerflowservice
          consumer:
            use-native-decoding: true
            # 显式指定分区数和主题一致
            partition-count: 2

3. 清理历史主题

之前错误配置下自动创建的pf-topic分区数不符合预期,需要先删除旧主题,避免配置不生效:

# 进入Kafka安装目录执行
bin/kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic pf-topic
验证方法
  1. 先启动Kafka集群,再启动2个消费端实例
  2. 启动生产者服务,调用/publish接口发送测试消息
  3. 查看消费端日志:node1键对应的a/b/c消息会全部落到一个分区,被其中一个实例消费;node2键对应的d/e/f消息会落到另一个分区,被第二个实例消费,符合预期。

补充:Kafka默认分区策略是对消息key的字节数组做Murmur2哈希后对分区数取模,你当前场景只有2个key、2个分区,两个key会恰好路由到不同分区,完全匹配你的并行消费预期。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 22:57:09