【发布时间】:2018-05-25 09:30:32
【问题描述】:
我正在使用 spring kafka (KafkaTemplate) 发送字符串消息。但为了使其与旧代码兼容,我需要在消息中附加额外的 CorrelationId。所以我创建了 ProducerRecord 对象,以我的消息作为它的值,并在它的 HeaderRecord 中设置 CorrelationId。
producerRecord = new ProducerRecord<>(
kafkaTemplate_.getDefaultTopic(),
null,
null,
null,
myStringMessage,
Collections.singletonList(new RecordHeader("CorrelationID", someIdAsBytes)));
kafkaTemplate_.sendDefault(producerRecord);
键和值序列化器设置为 StringSerializer,但 about 代码无法说明 ProduderRecord 不是 String 或 StringSerializer 类型。
如果我像下面那样执行 toString(),它会起作用。但是在MessageListener 端,接收到的ConsumerRecord 没有correlationId 作为它的RecordHeader,因为correlationId 是作为ProducerRecord 的RecordHeader 附加的。所以我必须在 MessageListener onMessage(Object msg) 上进行类型转换,例如从 Object 转换为 ConsumerRecord(这是 ProduderRecord 的字符串序列化版本),从 ConsumerRecord.value() 字段解析 ProducerRecord,然后从 ProcuderRecord 标头获取 CorrelationId,并从生产者记录。这看起来很麻烦。我的发送和接收逻辑正常吗?
kafkaTemplate_.sendDefault(producerRecord.toString());
【问题讨论】: