【问题标题】:How to set the kafka message key in the ProducerRecord如何在ProducerRecord中设置kafka message key
【发布时间】:2021-01-17 00:16:43
【问题描述】:

目前我正在使用 org.springframework.kafka.core.KafkaTemplate 发布带有标题的主题的 avro 消息。

@Override
public ListenableFuture<SendResult<K, V>> send(Message<?> message) {
    ProducerRecord<?, ?> producerRecord = this.messageConverter.fromMessage(message, this.defaultTopic);
    if (!producerRecord.headers().iterator().hasNext()) { // possibly no Jackson
        byte[] correlationId = message.getHeaders().get(KafkaHeaders.CORRELATION_ID, byte[].class);
        if (correlationId != null) {
            producerRecord.headers().add(KafkaHeaders.CORRELATION_ID, correlationId);
        }
    }
    return doSend((ProducerRecord<K, V>) producerRecord);
}

在 Message> 中,我们可以设置 value 和 headers 但不能设置 key。有没有办法在标题中输入密钥?如果是这样,请让我知道密钥的标题名称吗?有没有办法使用 KafkaTemplate 发送键、值和标头

【问题讨论】:

  • messageConverter.fromMessage 还没有设置 ProducerRecord 键? this.messageConverter 是什么类型的对象?

标签: spring-boot apache-kafka spring-kafka


【解决方案1】:

您必须在ProducerRecord 上设置标题,例如:

RecordHeader yourHeader = new RecordHeader("yourHeaderName", "yourValue".toByteArray())
record.headers().add(recordHeaderKafkaMessageKey)

编辑:误读要求,抱歉。消息密钥有一个定义的 kafka 标头,即KafkaHeaders.MESSAGE_KEY

【讨论】:

  • 实际上我需要填充消息键以及标题和值。我知道我可以设置标头,但发送标头不是问题。我无法使用 Message> 发送消息密钥
  • 对不起,应该仔细阅读您的问题。 KafkaHeaders.MESSAGE_KEY 有帮助吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-09-02
  • 2021-10-23
  • 1970-01-01
  • 2018-05-25
  • 2016-06-29
  • 2021-11-27
  • 2021-10-29
相关资源
最近更新 更多