You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 07:22:03