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

Kafka Streams在Strimzi集群中跨State Store查询返回Null求助

Kafka Streams双State Store在Strimzi集群上的查询异常问题

问题现象

  • 接收bp-addr记录时,可成功从bp-state-store获取数据,但从person-state-store查询始终返回Null;
  • 接收person-addr记录时,从bp-state-store查询也返回Null;
  • 使用.get()精确查询和.prefixScan()前缀扫描均出现该异常。

已完成的排查动作

  • 核对过键值、条目、State Store状态,均无异常;
  • 单元测试和本地Docker Kafka集群(cp-kafka)中运行正常;
  • 检查过Strimzi集群(kafka:0.29.0-kafka-3.1.0)的ACLs权限(已将用户设为超级用户),并在两个全新OpenShift命名空间中测试过,问题仍存在。

拓扑构建代码

private StreamsBuilder buildTopology(KafkaStreamsProperty kafkaStreamsProperty,
                                     SpecificAvroSerde<JoinedPersonAddrV2> joinedPASerde,
                                     SpecificAvroSerde<JoinedBpAddrV2> joinedBASerde,
                                     SpecificAvroSerde<ObjectUpdateEvent> updateSerde
) {
    StreamsBuilder builder = new StreamsBuilder();

    final Map<String, String> serdeConfig = Collections.singletonMap(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, kafkaStreamsProperty.getSchemaRegistryUrl());

    joinedBASerde.configure(serdeConfig, false);
    joinedPASerde.configure(serdeConfig, false);
    updateSerde.configure(serdeConfig, false);

    KStream<String, JoinedPersonAddrV2> personAddrKStream = builder.stream(kafkaStreamsProperty.getPersonInputTopic(), Consumed.with(Serdes.String(), joinedPASerde));
    KStream<String, JoinedBpAddrV2> bpAddrKStream = builder.stream(kafkaStreamsProperty.getBpInputTopic(), Consumed.with(Serdes.String(), joinedBASerde));

    StoreBuilder<KeyValueStore<String, JoinedBpAddrV2>> bpStoreBuilder = Stores.keyValueStoreBuilder(
            Stores.persistentKeyValueStore(BP_STORE_NAME),
            Serdes.String(),
            joinedBASerde
    );

    StoreBuilder<KeyValueStore<String, JoinedPersonAddrV2>> personStoreBuilder = Stores.keyValueStoreBuilder(
            Stores.persistentKeyValueStore(PERSON_STORE_NAME),
            Serdes.String(),
            joinedPASerde
    );

    StoreBuilder<KeyValueStore<String, Long>> debounceStoreBuilder = Stores.keyValueStoreBuilder(
            Stores.persistentKeyValueStore(DEBOUNCE_STORE),
            Serdes.String(),
            Serdes.Long()
    );

    builder.addStateStore(bpStoreBuilder);
    builder.addStateStore(personStoreBuilder);
    builder.addStateStore(debounceStoreBuilder);

    KStream<String, BpAddrOrPersonAddrV2> mergedStream = bpAddrKStream
            .mapValues(value -> new BpAddrOrPersonAddrV2(null, value))
            .merge(personAddrKStream.mapValues(value -> new BpAddrOrPersonAddrV2(value, null)));

    mergedStream.process(() -> new ContextualProcessor<String, BpAddrOrPersonAddrV2, String, ObjectUpdateEvent>() {
        @Override
        public void process(Record<String, BpAddrOrPersonAddrV2> record) {
            // assign value to a variable
            BpAddrOrPersonAddrV2 eitherValue = record.value();
            boolean isHashUpdated = false;

            // get state stores
            KeyValueStore<String, JoinedBpAddrV2> bpAddrStateStore = context().getStateStore(BP_STORE_NAME);
            KeyValueStore<String, JoinedPersonAddrV2> personAddrStateStore = context().getStateStore(PERSON_STORE_NAME);

            // assign single Values
            JoinedBpAddrV2 bpAddr = eitherValue.getBpAddr();
            JoinedPersonAddrV2 personAddr = eitherValue.getPersonAddr();

            // process if bp-addr record
            if (bpAddr != null) {
                String[] splittedKey = record.key().split(":");
                String personIdAddrIdKey = String.format("%s:%s", splittedKey[1], splittedKey[2]);

                boolean isEmittent = isEmittent(bpAddr);

                if (isRegisteredOwnerId(bpAddr, splittedKey[1]) || isEmittent) {
//                        personId, addrId, bpId
                    String revertedKey = String.format("%s:%s:%s", splittedKey[1], splittedKey[2], splittedKey[0]);

                    // get stored value from bp-addr state store
                    JoinedBpAddrV2 storedValue = bpAddrStateStore.get(revertedKey);
                    log.info("bpaddr entry storedValue: {}", storedValue);
                    if (storedValue == null) {
                        log.debug("{}, New BpAddr Table Entry for: {}, hash: {}", BP_STORE_NAME, record.key(), bpAddr.getHash());
                        bpAddrStateStore.put(revertedKey, bpAddr);
                        isHashUpdated = true;
                    } else if (!Objects.equals(storedValue.getHash(), bpAddr.getHash())) {
                        log.debug("{}, Update BpAddr Entry for key: {}, old hash: {}, new hash: {}", BP_STORE_NAME, record.key(), storedValue.getHash(), bpAddr.getHash());
                        bpAddrStateStore.put(revertedKey, bpAddr);
                        isHashUpdated = true;
                    } else {
                        log.debug("{}, No hash change do nothing for BpAddr key:{} hash value:{}", BP_STORE_NAME, record.key(), storedValue.getHash());
                    }
                }

                if (isHashUpdated) {
                    log.info("SEARCHING IN PERSON - ADDR STATE STORE WITH A KEY: {}", personIdAddrIdKey);

                    JoinedPersonAddrV2 matching = personAddrStateStore.get(personIdAddrIdKey);
                    log.info("MATCHING: {}", matching);

                    KeyValueIterator<String, JoinedPersonAddrV2> matchedPersonAddrIterator = personAddrStateStore.prefixScan(personIdAddrIdKey, new StringSerializer());

                    if (matchedPersonAddrIterator.hasNext()) {
                        while (matchedPersonAddrIterator.hasNext()) {
                            JoinedPersonAddrV2 matchedPersonAddr = matchedPersonAddrIterator.next().value;
                            // if the bp-addr is an emmitent and didn't find mathing personAddr look with just personId using prefix scan
                            // because addrId is not propagated on bp, therefore they might have a different domiAddr on bp and person
                            if (matchedPersonAddr == null && isEmittent) {
                                String personId = splittedKey[1];
                                // do a prefix scan for a personId
                                KeyValueIterator<String, JoinedPersonAddrV2> joinedPersonIterator = personAddrStateStore.prefixScan(personId, new StringSerializer());
                                while (joinedPersonIterator.hasNext()) {
                                    JoinedPersonAddrV2 item = joinedPersonIterator.next().value;
                                    context().forward(record.withKey(personIdAddrIdKey).withValue(createObjectUpdateEvent(item, bpAddr)));
                                }
                                // "standard" match forward the record
                            } else if (matchedPersonAddr != null) {
                                context().forward(record.withKey(personIdAddrIdKey).withValue(createObjectUpdateEvent(matchedPersonAddr, bpAddr)));
                            }
                        }
                    }
                }

            } else if (personAddr != null) { // process if person-addr record is coming

                if (record.value() != null) {
                    JoinedPersonAddrV2 storedValue = personAddrStateStore.get(record.key());
                    log.info("personaddr entry storedValue: {}", storedValue);

                    if (storedValue == null) {
                        log.debug("{}, New PersonAddr Table entry for: {}, hash: {}", PERSON_STORE_NAME, record.key(), personAddr.getHash());
                        personAddrStateStore.put(record.key(), personAddr);
                        isHashUpdated = true;
                    } else if (!storedValue.getHash().equals(personAddr.getHash())) {
                        log.debug("{}, Update PersonAddr Table entry for: {}, old hash: {}, new hash: {}", PERSON_STORE_NAME, record.key(), storedValue.getHash(), personAddr.getHash());
                        personAddrStateStore.put(record.key(), personAddr);
                        isHashUpdated = true;
                    } else {
                        log.debug("{}, No hash change do nothing for PersonAddr key: {}, hash: {}", PERSON_STORE_NAME, record.key(), storedValue.getHash());
                    }
                } else {
                    log.debug("Got NULL JoinedPersonAddr Object - skip! --> key: {}", record.key());
                }

                if (isHashUpdated) {
                    try {

                        // look for matches in bpAddr store
                        log.info("SEARCHING IN BP - ADDR STATE STORE WITH A KEY: {}", record.key());
                        KeyValueIterator<String, JoinedBpAddrV2> bpAddrIterator = bpAddrStateStore.prefixScan(record.key(), new StringSerializer());

                        if (bpAddrIterator.hasNext()) {
                            log.trace("Found records for key {} record in {}", record.key(), BP_STORE_NAME);
                            // got matches for given personId:addrId update all found records
                            while (bpAddrIterator.hasNext()) {
                                JoinedBpAddrV2 bpAddrObject = bpAddrIterator.next().value;
                                context().forward(record.withKey(record.key()).withValue(createObjectUpdateEvent(personAddr, bpAddrObject)));
                            }
                        } else {
                            log.trace("Found NO records for key {} record in {}", record.key(), BP_STORE_NAME);
                            String[] splittedKey = record.key().split(":");
                            final String ZERO = "0";

                            // only when addr != 0 and is emittent
                            if (personAddr.getIsEmittent() && !splittedKey[1].equals(ZERO)) {
                                String personIdNullAddrKey = String.format("%s:%s", splittedKey[0], ZERO);
//                                look for records in bp-addr state store with key personId:0
                                KeyValueIterator<String, JoinedBpAddrV2> bpNullAddrIterator = bpAddrStateStore.prefixScan(personIdNullAddrKey, new StringSerializer());
                                if (bpNullAddrIterator.hasNext()) {
                                    log.trace("Found records for key {} record in {}", personIdNullAddrKey, BP_STORE_NAME);
                                    while (bpNullAddrIterator.hasNext()) {
                                        JoinedBpAddrV2 bpAddrV2 = bpNullAddrIterator.next().value;
                                        context().forward(record.withKey(record.key()).withValue(createObjectUpdateEvent(personAddr, bpAddrV2)));
                                    }
                                }
                            } else {
                                log.trace("Found NO matched sending PERSON_ONLY records");
                                // no match in bpAddr state store send PERSON_ONLY event
                                context().forward(record.withValue(createObjectUpdateEvent(personAddr, null)));
                            }
                        }
                    } catch (Exception e) {
                        log.error("EXCEPTION : {}", e.getMessage());
                    }
                }
            }
        }
    }, BP_STORE_NAME, PERSON_STORE_NAME)
            .process(() -> new DebounceTransformer<>(DEBOUNCE_STORE, 5000), DEBOUNCE_STORE)
            .peek((key, value) -> log.debug("Produced Update Event key: {}, hash: {}", key, value.getHash()))
            .to(kafkaStreamsProperty.getOutputTopic(), Produced.with(Serdes.String(), updateSerde));

    return builder;
}

拓扑结构

Topology: Topologies:
   Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000 (topics: [person.addr.join.topic])
      --> KSTREAM-MAPVALUES-0000000003
    Source: KSTREAM-SOURCE-0000000001 (topics: [bp.addr.join.topic])
      --> KSTREAM-MAPVALUES-0000000002
    Processor: KSTREAM-MAPVALUES-0000000002 (stores: [])
      --> KSTREAM-MERGE-0000000004
      <-- KSTREAM-SOURCE-0000000001
    Processor: KSTREAM-MAPVALUES-0000000003 (stores: [])
      --> KSTREAM-MERGE-0000000004
      <-- KSTREAM-SOURCE-0000000000
    Processor: KSTREAM-MERGE-0000000004 (stores: [])
      --> KSTREAM-PROCESSOR-0000000005
      <-- KSTREAM-MAPVALUES-0000000002, KSTREAM-MAPVALUES-0000000003
    Processor: KSTREAM-PROCESSOR-0000000005 (stores: [bp-state-store, person-state-store])
      --> KSTREAM-PROCESSOR-0000000006
      <-- KSTREAM-MERGE-0000000004
    Processor: KSTREAM-PROCESSOR-0000000006 (stores: [debounce-store])
      --> KSTREAM-PEEK-0000000007
      <-- KSTREAM-PROCESSOR-0000000005
    Processor: KSTREAM-PEEK-0000000007 (stores: [])
      --> KSTREAM-SINK-0000000008
      <-- KSTREAM-PROCESSOR-0000000006
    Sink: KSTREAM-SINK-0000000008 (topic: update.event.topic)
      <-- KSTREAM-PEEK-0000000007

排查思路与建议

  1. 统一键序列化逻辑
    代码中prefixScan使用new StringSerializer(),但State Store定义时用的是Serdes.String(),虽然都是String序列化,但可能存在编码或配置差异。建议替换为Serdes.String().serializer(),确保存储和查询时的键序列化逻辑完全一致。

  2. 验证流的分区一致性
    Kafka Streams的State Store是分区级别的,只有当两个输入主题的分区数相同,且键的分区策略一致时,同一键的数据才会落在同一个处理器实例的State Store分区中。检查person.addr.join.topic和bp.addr.join.topic的分区数,以及生产者是否使用相同的分区器(默认按键哈希,自定义分区器需确保逻辑一致)。

  3. 检查State Store的持久化与加载状态

    • 确认Strimzi环境下state.dir目录的权限和可用空间,避免磁盘配额或权限问题导致数据未正确落地;
    • 应用重启后,检查State Store的预加载是否完成,可通过KafkaStreams#store()方法手动查询键值,确认数据是否正确加载。
  4. 启用调试日志定位细节
    开启org.apache.kafka.streams和org.apache.kafka.streams.state的DEBUG日志,重点关注:

    • State Store的读写操作日志,查看键的实际序列化后的值;
    • 分区分配日志,确认处理器实例是否分配到了对应的数据分区;
    • Avro序列化/反序列化日志,排查是否存在Schema兼容性问题导致数据存储异常。
  5. 检查Avro Schema的集群兼容性
    确认集群Schema Registry中JoinedPersonAddrV2和JoinedBpAddrV2的Schema版本与本地一致,避免因Schema字段变更导致反序列化后的数据无法正确匹配查询键。

内容的提问来源于stack exchange,提问作者bohunn

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 06:00:53