【发布时间】:2021-06-11 21:30:33
【问题描述】:
我的消费者配置有 kafka 批处理侦听器配置和 @KafkaListener 消费消息列表。我有一个 ConsumerInterceptor,我想为每条记录设置唯一的 id,并将其值存储在映射诊断上下文 (MDC) 中。如果我的 kafka 侦听器使用单个消息,则唯一 ID 是正确的。但是我的 kafka 侦听器使用消息列表,因此 MDC.get("id") 仅获得最后一个值。怎么可能处理? 我的拦截器;
public class KafkaConsumerInterceptor implements ConsumerInterceptor<String, String>{
@Override
public ConsumerRecords<String, String> onConsume(ConsumerRecords<String, String> consumerRecords) {
ConsumerRecord<String, String> record = consumerRecords.iterator().next();
setId(record.headers().headers("id"));
return consumerRecords;
}
private void setId(Iterable<Header> idHeader) {
String id = UUID.randomUUID().toString();
if (idHeader.iterator().hasNext()) {
Header header = idHeader.iterator().next();
id = new String(header.value(), StandardCharsets.UTF_8);
}
MDC.put("id", id);
}
}
【问题讨论】:
标签: apache-kafka spring-kafka mdc