【问题标题】:How to read Kafka Message Key from Spring cloud streams?如何从 Spring 云流中读取 Kafka 消息密钥?
【发布时间】:2019-07-06 17:01:59
【问题描述】:

我正在使用 Spring Cloud Streams 来消费来自 Kafka 的消息。

是否可以从代码中读取Kafka Message Key?

我有一个 Kafka 主题,通常有 2 种类型的消息。要采取的操作因消息键而异。我看到 spring 文档只有以下内容来阅读消息。在这里,我需要指定消息的实际映射(此处为 Greetings 类)。但是,我需要一种方法来读取消息密钥并确定可反序列化的 Pojo

public class GreetingsListener {

    @StreamListener(GreetingsProcessor.INPUT)
    public void handleGreetings(@Payload Greetings request) {
     
    }
}

【问题讨论】:

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


    【解决方案1】:

    你可以试试这样的:

    @StreamListener(GreetingsProcessor.INPUT)
    public void handleGreetings(@Payload Greetings request, @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY)String  key) {
    
    }
    

    您需要为密钥提供适当的反序列化程序。例如如果您的密钥是字符串,那么您可以提供:

    spring.cloud.stream.kafka.binder.configuration.key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
    

    如果需要为不同的输入通道使用不同的密钥解串器,可以在每个 kafka 绑定的 producer 部分下扩展此设置。例如:

    spring:
      cloud:
        stream:
          kafka:
            bindings:
              <channel_name>:
                consumer:
                  startOffset: latest
                  autoCommitOffset: true
                  autoCommitOnError: true
                  configuration:
                    key.deserializer: org.apache.kafka.common.serialization.StringDeserializer
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-09-16
      • 1970-01-01
      • 1970-01-01
      • 2020-04-30
      • 2019-09-20
      • 1970-01-01
      相关资源
      最近更新 更多