【问题标题】:KafkaHeaders.RECEIVED_MESSAGE_KEY vs KafkaHeaders.MESSAGE_KEY header on Spring CloudSpring Cloud 上的 KafkaHeaders.RECEIVED_MESSAGE_KEY 与 KafkaHeaders.MESSAGE_KEY 标头
【发布时间】:2020-03-23 16:43:00
【问题描述】:

我正在使用弹簧云流。 我想知道KafkaHeaders.RECEIVED_MESSAGE_KEYKafkaHeaders.MESSAGE_KEY 有什么区别

我有 2 个项目,第一个使用 KafkaHeaders.MESSAGE_KEY 作为标头生成消息:

    public void sendResponse(ThirdPartyResponse thirdPartyResponse) {

        log.info("Sending response of type 'completed' [{}].", thirdPartyResponse);
        integrations.send(
                withPayload(ApplicationSubmissionSuccessPayload.success(thirdPartyResponse))
                        .setHeader(KafkaHeaders.MESSAGE_KEY, thirdPartyResponse.getData().getApplicationId())
                        .build());

    }

第二个使用KafkaHeaders.RECEIVED_MESSAGE_KEY消费

@StreamListener(target = "ofaOut")
public void receive(@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) String applicationId, @Payload String payload) throws JsonProcessingException {

...
}

但是我得到了这个错误

    2020-03-23 16:13:27.924 ERROR 1 --- [container-0-C-1] o.s.integration.handler.LoggingHandler   : 
org.springframework.messaging.MessageHandlingException: Missing header 'kafka_receivedMessageKey' for method parameter type [class java.lang.String], failedMessage=GenericMessage [payload=byte[739],
 headers={kafka_offset=285, scst_nativeHeadersPresent=true, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@67c19b7c, deliveryAttempt=3, kafka_timestampType=CREATE_TIME, 
kafka_receivedMessageKey=null, kafka_receivedPartitionId=0, 
contentType=application/json, kafka_receivedTopic=com.product.foo.ofa.out, kafka_receivedTimestamp=1584715870225, kafka_groupId=aop-foo-kyc}]

它缺少标题

Missing header 'kafka_receivedMessageKey'

我该如何解决?

【问题讨论】:

    标签: java spring-cloud spring-cloud-stream


    【解决方案1】:

    RECEIVED... 设置在入站消息上;另一种是让应用程序指定出站消息的键值。

    当应用程序接收到消息并执行某些工作并将消息重新发布到另一个主题时,它们是不同的以避免意外传播。

    使用 Spring Integration 时,标头会随着消息遍历流而自动复制。

    出站消息映射器不会映射RECEIVED... 标头,因此它们不会出现在ProducerRecord 中。

    ... kafka_receivedMessageKey=null ...

    表示入站记录中的键为空。

    要接收空键,请使用

    @Header(name = KafkaHeaders.RECEIVED_MESSAGE_KEY, required = false)
    

    【讨论】:

    • 所以我应该在我的生产者项目中做.setHeader(KafkaHeaders. RECEIVED_MESSAGE_KEY, thirdPartyResponse.getData().getApplicationId()) ??
    • 否;您所拥有的是正确的-只是您在该主题中有一个没有密钥的记录;要允许接收此类消息,您需要将 required=false 添加到注释中。
    • 知道了,但是如果我在做.setHeader(KafkaHeaders.MESSAGE_KEY, thirdPartyResponse.getData().getApplicationId()),为什么我要发送带有空键的消息@
    • 可能是您开始设置密钥之前的“旧”记录?或者getApplicationId() 返回了null?您可以使用控制台使用者$ kafka-console-consumer --bootstrap-server localhost:9092 --topic myTopic --from-beginning --property key.separator=: --property print.key=true 检查记录。它将打印为key:value
    猜你喜欢
    • 2021-06-04
    • 2016-07-24
    • 2022-07-01
    • 1970-01-01
    • 2019-02-02
    • 2021-02-28
    • 1970-01-01
    • 1970-01-01
    • 2016-03-29
    相关资源
    最近更新 更多