【问题标题】:Spring Cloud Stream send/consume message to different partitions with KafkaHeaders.Message_KEYSpring Cloud Stream 使用 KafkaHeaders.Message_KEY 向不同分区发送/消费消息
【发布时间】:2022-07-01 06:39:41
【问题描述】:
I am trying to implement a prototype for implementing messaging system using Spring Cloud Stream. I selected Apache Kafka as binder. I created a topic with 2 partitions for scalability.  Then I tried to  send different messages to different partitions using following rest api method. 

我为 2 个分区设置了 2 个不同的消息键。

@PostMapping(\"/publish\")
public void publish(@RequestParam String message) {
    log.debug(\"REST request the message : {} to send to Kafka topic \", message);
    Message message1 = MessageBuilder.withPayload(\"Hello from a\")
        .setHeader(KafkaHeaders.MESSAGE_KEY, \"node1\")
        .build();
    Message message2 = MessageBuilder.withPayload(\"Hello from b\")
        .setHeader(KafkaHeaders.MESSAGE_KEY, \"node1\")
        .build();
    Message message3 = MessageBuilder.withPayload(\"Hello from c\")
        .setHeader(KafkaHeaders.MESSAGE_KEY, \"node1\")
        .build();
    Message message4 = MessageBuilder.withPayload(\"Hello from d\")
        .setHeader(KafkaHeaders.MESSAGE_KEY, \"node2\")
        .build();
    Message message5 = MessageBuilder.withPayload(\"Hello from e\")
        .setHeader(KafkaHeaders.MESSAGE_KEY, \"node2\")
        .build();
    Message message6 = MessageBuilder.withPayload(\"Hello from f\")
        .setHeader(KafkaHeaders.MESSAGE_KEY, \"node2\")
        .build();
    output.send(\"simulatePf-out-0\", message1);
    output.send(\"simulatePf-out-0\", message2);
    output.send(\"simulatePf-out-0\", message3);
    output.send(\"simulatePf-out-0\", message4);
    output.send(\"simulatePf-out-0\", message5);
    output.send(\"simulatePf-out-0\", message6);


}

这是我用于生产者应用程序的 application.yml

  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
  bindings:
    simulatePf-out-0:
      producer:
        useNativeEncoding: true
        partition-count: 3
      destination: pf-topic
      content-type: text/plain
      group: dsa-back-end

为了测试并行性,我创建了一个从 pf-topic 读取消息的消费者应用程序。这是来自消费者应用程序的配置。

  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

  bindings:
    simulatePf-in-0:
      destination: pf-topic
      content-type: text/plain
      group: powerflowservice
      consumer:
        use-native-decoding: true

. 我在消费者应用程序中创建了一个函数来消费消息

   @Bean
public Consumer<Message> simulatePf() {
    return message -> {
        log.info(\"header \" + message.getHeaders());
        log.info(\"received \" + message.getPayload());
    };
}

现在是测试的时候了。为了测试并行性,我运行了 2 个 spring boot consumer application 实例。我期待看到一个消费者从一个分区消费消息,其他消费者消费者从另一个分区消费消息。所以我希望消息a,消息b,消息被消费者一消费。消息 d、消息 e 和消息 f 是其他消费者的消费者。因为我设置了不同的消息键来分配不同的分区。但所有消息仅由一个应用程序使用

 2022-06-30 20:34:48.895  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : header {deliveryAttempt=1, kafka_timestampType=CREATE_TIME, kafka_receivedMessageKey=node1, kafka_receivedTopic=pf-topic, skip-input-type-conversion=true, kafka_offset=270, scst_nativeHeadersPresent=true, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@1eaf51df, source-type=streamBridge, id=a77d12f2-f184-0f2f-6a76-147803dd43f3, kafka_receivedPartitionId=0, kafka_receivedTimestamp=1656610488838, kafka_groupId=powerflowservice, timestamp=1656610488890}
2022-06-30 20:34:48.901  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : received Hello from a
2022-06-30 20:34:48.929  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : header {deliveryAttempt=1, kafka_timestampType=CREATE_TIME, kafka_receivedMessageKey=node1, kafka_receivedTopic=pf-topic, skip-input-type-conversion=true, kafka_offset=271, scst_nativeHeadersPresent=true, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@1eaf51df, source-type=streamBridge, id=2e89f9b7-b6e7-482f-3c46-f73b2ad0705c, kafka_receivedPartitionId=0, kafka_receivedTimestamp=1656610488840, kafka_groupId=powerflowservice, timestamp=1656610488929}
2022-06-30 20:34:48.932  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : received Hello from b
2022-06-30 20:34:48.933  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : header {deliveryAttempt=1, kafka_timestampType=CREATE_TIME, kafka_receivedMessageKey=node1, kafka_receivedTopic=pf-topic, skip-input-type-conversion=true, kafka_offset=272, scst_nativeHeadersPresent=true, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@1eaf51df, source-type=streamBridge, id=15640532-b57f-b58e-62e7-c2bc9375fdf0, kafka_receivedPartitionId=0, kafka_receivedTimestamp=1656610488841, kafka_groupId=powerflowservice, timestamp=1656610488933}
2022-06-30 20:34:48.934  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : received Hello from c
2022-06-30 20:34:48.935  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : header {deliveryAttempt=1, kafka_timestampType=CREATE_TIME, kafka_receivedMessageKey=node2, kafka_receivedTopic=pf-topic, skip-input-type-conversion=true, kafka_offset=273, scst_nativeHeadersPresent=true, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@1eaf51df, source-type=streamBridge, id=590f0fb7-042f-e134-d214-ead570e42fe3, kafka_receivedPartitionId=0, kafka_receivedTimestamp=1656610488842, kafka_groupId=powerflowservice, timestamp=1656610488934}
2022-06-30 20:34:48.938  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : received Hello from d
2022-06-30 20:34:48.940  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : header {deliveryAttempt=1, kafka_timestampType=CREATE_TIME, kafka_receivedMessageKey=node2, kafka_receivedTopic=pf-topic, skip-input-type-conversion=true, kafka_offset=274, scst_nativeHeadersPresent=true, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@1eaf51df, source-type=streamBridge, id=9a67e68b-95d4-a02e-cc14-ac30c684b639, kafka_receivedPartitionId=0, kafka_receivedTimestamp=1656610488842, kafka_groupId=powerflowservice, timestamp=1656610488940}
2022-06-30 20:34:48.941  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : received Hello from e
2022-06-30 20:34:48.943  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : header {deliveryAttempt=1, kafka_timestampType=CREATE_TIME, kafka_receivedMessageKey=node2, kafka_receivedTopic=pf-topic, skip-input-type-conversion=true, kafka_offset=275, scst_nativeHeadersPresent=true, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@1eaf51df, source-type=streamBridge, id=333269af-bbd5-12b0-09de-8bd7959ebf08, kafka_receivedPartitionId=0, kafka_receivedTimestamp=1656610488843, kafka_groupId=powerflowservice, timestamp=1656610488943}
2022-06-30 20:34:48.943  INFO 11860 --- [container-0-C-1] c.s.powerflow.config.AsyncConfiguration  : received Hello from f

你能帮我解决我所缺少的吗?

  • 在我看来,您的主题只有一个分区,或者您的密钥都从生产者散列到同一个分区中。您可以使用kafka-consumer-groups --describe 查看为哪些消费者分配了哪些分区
  • 不。它有 3 个分区 主题:pf-topic 分区:0 领导者:3 副本:3,1 Isr:3,1 主题:pf-topic 分区:1 领导者:1 副本:1,2 Isr:1,2 主题: pf-topic 分区:2 领导者:2 副本:2,3 Isr:2,3
  • 好吧,好吧,仍然有可能两个键可以具有相同的哈希,所以最终在同一个分区中

标签: java apache-kafka jhipster spring-cloud-stream


【解决方案1】:

您仅在发送时将消息键设置为标头。您可以在消息上添加KafkaHeaders.PARTITION 标头以强制执行特定分区。

如果您不想通过标头添加硬编码分区,则可以在应用程序中设置分区键 SpEL 表达式或分区键提取器 bean。这两种机制都是 Spring Cloud Stream 特有的。如果您提供其中任何一个,您仍然需要告诉 Spring Cloud Stream 您要如何选择分区。为此,您可以使用分区选择器 SpEL 表达式或分区选择器策略。如果您不提供它们,那么它将使用默认选择器策略,方法是采用消息键 % 主题分区数的 hashCode

我想你昨天问了另一个相关问题,我在回答中链接了这个blog。在该博客的最后几节中,解释了所有这些细节。

引用博客:

如果您不提供分区键表达式或分区键提取器 bean,那么 Spring Cloud Stream 将完全不为您做出任何分区决策。在这种情况下,如果主题有多个分区,则会触发 Kafka 的默认分区机制。默认情况下,Kafka 使用 DefaultPartitioner,如果消息有一个键(见上文),则使用该键的哈希来计算分区。

我认为您在应用程序中看到了 Kafka 的默认行为。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-06-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-06-21
    • 2021-04-27
    • 2016-09-09
    • 2018-03-09
    相关资源
    最近更新 更多