如何在Kafka Stream的aggregate()函数中读取Record Headers?
在Kafka Streams的aggregate()中读取Record Headers
场景描述
发送Kafka消息时会携带用于标识消息类型、辅助业务处理的Header,示例代码如下:
ProducerRecord<String, String> producerRecord = new ProducerRecord<>(topic, key, msg); producerRecord.headers().add("headerkey", type.getBytes(StandardCharsets.UTF_8)); producer.send(producerRecord);
之前可以通过Processor API读取Header:
streamsBuilder.stream("topic", Consumed.with(Serdes.String(), Serdes.String())) .process(() -> new Processor(), "stateStore"); // 可在Processor类中读取Header
但业务中使用aggregate()做聚合处理时,需要在聚合逻辑里直接用到Header的值,当前的聚合代码示例如下:
stream.aggregate(String::new, (key, value, aggregated) -> aggregateValues(value, aggregated), Materialized.with(Serdes.String(), Serdes.String())) // 此处需要Header值
解决方案
直接在aggregate()的参数里无法获取Header,需要先通过transformValues将Header与原消息值封装到同一个对象中,再对转换后的流做聚合操作。
1. 定义封装类
创建一个包含原消息值和Header信息的类,用于传递数据:
public class ValueWithHeader { private final String value; private final String headerValue; public ValueWithHeader(String value, String headerValue) { this.value = value; this.headerValue = headerValue; } public String getValue() { return value; } public String getHeaderValue() { return headerValue; } }
2. 转换流并提取Header
使用transformValues处理原流,从ProcessorContext中获取当前Record的Header,封装到自定义对象中:
StreamsBuilder streamsBuilder = new StreamsBuilder(); KStream<String, ValueWithHeader> streamWithHeader = streamsBuilder.stream("topic", Consumed.with(Serdes.String(), Serdes.String())) .transformValues(() -> new ValueTransformer<String, ValueWithHeader>() { private ProcessorContext context; @Override public void init(ProcessorContext context) { this.context = context; } @Override public ValueWithHeader transform(String value) { // 提取指定key的Header值 String headerValue = null; Headers headers = context.headers(); for (Header header : headers) { if ("headerkey".equals(header.key())) { headerValue = new String(header.value(), StandardCharsets.UTF_8); break; } } return new ValueWithHeader(value, headerValue); } @Override public void close() { // 按需清理资源 } });
3. 在aggregate中使用Header值
现在对转换后的流执行聚合,就能在逻辑中直接获取Header的值了:
streamWithHeader.aggregate( String::new, // 初始化聚合结果 (key, valueWithHeader, aggregated) -> { String originalValue = valueWithHeader.getValue(); String headerValue = valueWithHeader.getHeaderValue(); // 结合Header值执行你的聚合逻辑 return aggregateValues(originalValue, aggregated, headerValue); }, Materialized.with(Serdes.String(), Serdes.String()) );
补充说明
如果使用Kafka Streams 3.0及以上版本,也可以利用Record相关API简化操作,但上述方法兼容绝大多数版本,适用性更广。
内容的提问来源于stack exchange,提问作者Samy
相关产品推荐
相关产品推荐

