如何在Kafka Streams的GlobalKTable中访问Kafka消息头
解决GlobalKTable存储并访问Kafka消息头的方案
默认情况下,Kafka Streams的GlobalKTable只存储消息的Key和Value,不会保留Headers。要实现需求,需要通过自定义包装类+转换逻辑将Headers与Value一起存入GlobalKTable,同时调整KStream的启动时机确保GlobalKTable加载完成。
1. 定义包装类封装Value与Headers
创建一个类将GenericRecord和Kafka Headers打包,方便存入状态存储:
import org.apache.kafka.common.header.Headers; import org.apache.kafka.common.header.internals.RecordHeaders; import org.apache.avro.generic.GenericRecord; public class RecordWithHeaders { private GenericRecord value; private Headers headers; public RecordWithHeaders(GenericRecord value, Headers headers) { this.value = value; this.headers = new RecordHeaders(headers); // 复制Headers避免引用冲突 } // Getter方法 public GenericRecord getValue() { return value; } public Headers getHeaders() { return headers; } }
2. 实现包装类的Serde
由于Kafka的Headers类无法直接序列化,需要自定义Serde来处理RecordWithHeaders的序列化与反序列化:
import org.apache.kafka.common.serialization.Deserializer; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serializer; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.module.SimpleModule; import org.apache.kafka.common.header.Header; import org.apache.kafka.common.header.internals.RecordHeaders; import java.util.Map; public class RecordWithHeadersSerde implements Serde<RecordWithHeaders> { private final ObjectMapper objectMapper; public RecordWithHeadersSerde() { this.objectMapper = new ObjectMapper(); SimpleModule module = new SimpleModule(); // 序列化Headers为字节数组 module.addSerializer(Headers.class, (Serializer<Headers>) (topic, data) -> { if (data == null) return null; return objectMapper.writeValueAsBytes(data.toArray()); }); // 反序列化字节数组为Headers module.addDeserializer(Headers.class, (Deserializer<Headers>) (topic, data) -> { if (data == null) return null; Header[] headers = objectMapper.readValue(data, Header[].class); RecordHeaders recordHeaders = new RecordHeaders(); for (Header header : headers) { recordHeaders.add(header.key(), header.value()); } return recordHeaders; }); objectMapper.registerModule(module); } @Override public void configure(Map<String, ?> configs, boolean isKey) {} @Override public void close() {} @Override public Serializer<RecordWithHeaders> serializer() { return (topic, data) -> { if (data == null) return null; try { return objectMapper.writeValueAsBytes(data); } catch (Exception e) { throw new RuntimeException("序列化RecordWithHeaders失败", e); } }; } @Override public Deserializer<RecordWithHeaders> deserializer() { return (topic, data) -> { if (data == null) return null; try { return objectMapper.readValue(data, RecordWithHeaders.class); } catch (Exception e) { throw new RuntimeException("反序列化RecordWithHeaders失败", e); } }; } }
3. 修改GlobalKTable配置,拦截记录并包装
调整GlobalKTable的创建逻辑,在消费消息时将Headers与Value打包存入状态存储:
@Bean(name = "createGlobalKTable") public GlobalKTable<String, RecordWithHeaders> createGlobalKTable(@Qualifier("globalKTableStreamsBuilder") StreamsBuilderFactoryBean globalKTableStreamsBuilderFactoryBean) throws Exception { StreamsBuilder streamsBuilder = globalKTableStreamsBuilderFactoryBean.getObject(); Map<String, String> serdeConfig = new HashMap<>(); serdeConfig.put("schema.registry.url", schemaRegistryUrl); serdeConfig.put("schema.registry.basic.auth.user.info", basicAuthUserInfo); serdeConfig.put("schema.registry.basic.auth.credentials.source", basicAuthCredentialsSource); final Serde<GenericRecord> valueGenericAvroSerde = new GenericAvroSerde(); valueGenericAvroSerde.configure(serdeConfig, false); // 配置消费时的Serde Consumed<String, GenericRecord> consumed = Consumed.with(Serdes.String(), valueGenericAvroSerde); return streamsBuilder.globalTable(ediLegTopic, consumed, Materialized.<String, RecordWithHeaders, KeyValueStore<Bytes, byte[]>>as(storeName + "-global") .withKeySerde(Serdes.String()) .withValueSerde(new RecordWithHeadersSerde())) // 转换Value为包含Headers的包装类 .transformValues(() -> new ValueTransformerWithKey<String, GenericRecord, RecordWithHeaders>() { private ProcessorContext context; @Override public void init(ProcessorContext context) { this.context = context; } @Override public RecordWithHeaders transform(String key, GenericRecord value) { // 从上下文获取当前消息的Headers Headers headers = context.headers(); return new RecordWithHeaders(value, headers); } @Override public void close() {} }); }
4. 调整GlobalKTableService读取逻辑
修改服务类,从GlobalKTable中读取包装类以获取Headers:
@Service public class GlobalKTableService { private final String storeName; private final KafkaStreams globalKafkaStreams; public GlobalKTableService(@Value("${store.name}") String storeName, @Qualifier("globalKTableStreamsBuilder") StreamsBuilderFactoryBean globalStreamsFactory) { this.storeName = storeName + "-global"; this.globalKafkaStreams = globalStreamsFactory.getKafkaStreams(); } public RecordWithHeaders getValueFromGlobalKTable(String key) { ReadOnlyKeyValueStore<String, RecordWithHeaders> store = globalKafkaStreams.store( StoreQueryParameters.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore()) ); return store.get(key); } public KeyValueIterator<String, RecordWithHeaders> getAllValuesFromGlobalKTable() { ReadOnlyKeyValueStore<String, RecordWithHeaders> store = globalKafkaStreams.store( StoreQueryParameters.fromNameAndType(storeName, QueryableStoreTypes.keyValueStore()) ); return store.all(); } }
5. 更新CustomProcessor处理逻辑
在处理器中读取包装类,获取GlobalKTable中存储的Headers:
@Override public void process(Record<String, GenericRecord> processingRecord) { String key = processingRecord.key(); GenericRecord streamValue = processingRecord.value(); Headers streamHeaders = processingRecord.headers(); RecordWithHeaders globalRecord = globalKTableService.getValueFromGlobalKTable(key); if (globalRecord == null) { log.info("GlobalKTable中无key: {}的记录", key); logicProvider.applyLogic(new KeyValue<>(key, streamValue), streamHeaders); context.forward(processingRecord.withValue(streamValue)); return; } GenericRecord globalValue = globalRecord.getValue(); Headers globalHeaders = globalRecord.getHeaders(); if (!streamValue.equals(globalValue)) { log.info("GlobalKTable中无匹配记录,转发key: {}的消息", key); // 此处可使用globalHeaders进行业务处理 logicProvider.applyLogic(new KeyValue<>(key, streamValue), streamHeaders); context.forward(processingRecord.withValue(streamValue)); } else { log.info("GlobalKTable中找到匹配记录,跳过key: {}的消息", key); } }
6. 确保GlobalKTable加载完成后启动KStream
将KStream的自动启动关闭,等待GlobalKTable进入RUNNING状态后再启动:
// GlobalKTable的Builder配置 @Bean(name = "globalKTableStreamsBuilder") public StreamsBuilderFactoryBean globalKTableStreamsBuilder(KafkaStreamsConfiguration globalKTableConfig) { StreamsBuilderFactoryBean factoryBean = new StreamsBuilderFactoryBean(globalKTableConfig); factoryBean.setKafkaStreamsCustomizer(kafkaStreams -> { kafkaStreams.setStateListener((newState, oldState) -> { log.info("GlobalKTable状态变更: {} -> {}", oldState, newState); if (newState == KafkaStreams.State.RUNNING) { // GlobalKTable加载完成,启动KStream kStreamBuilder.getKafkaStreams().start(); } }); }); factoryBean.setAutoStartup(true); // 自动启动GlobalKTable return factoryBean; } // KStream的Builder配置 @Bean public StreamsBuilderFactoryBean kStreamBuilder( @Qualifier("kafkaStreamsConfig") KafkaStreamsConfiguration config, GlobalKTableService globalKTableService) { StreamsBuilderFactoryBean factoryBean = new StreamsBuilderFactoryBean(config, new CleanupConfig(false, false)); factoryBean.setAutoStartup(false); // 禁止自动启动 factoryBean.setKafkaStreamsCustomizer(kafkaStreams -> { kafkaStreams.setStateListener((newState, oldState) -> { log.info("KStream状态变更: {} -> {}", oldState, newState); if (newState == KafkaStreams.State.RUNNING) { KeyValueIterator<String, RecordWithHeaders> iterator = globalKTableService.getAllValuesFromGlobalKTable(); // 执行初始化逻辑 } }); }); return factoryBean; }
内容的提问来源于stack exchange,提问作者Misa
相关产品推荐
相关产品推荐

