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

Kafka Streams处理器无法访问StateStore问题求助

Kafka Streams StateStore 无法关联处理器问题

问题描述

编写了代码A,调用builtTopology = builder.build后,执行new org.apache.kafka.streams.TopologyTestDriver(builtTopology, properties)时出现错误B。已完成以下检查:

  • Builder配置:通过builder.addStateStore注册StateStore,并在调用process方法时传入存储名称显式关联处理器;
  • ProcessorContext与StateStore:在处理器的init初始化阶段,通过processor context和存储名称获取了StateStore,确保处理器可访问该存储。

仍无法解决问题,请问遗漏了什么?

错误信息B

Caused by: org.apache.kafka.streams.errors.StreamsException: Processor KSTREAM-PROCESSOR-0000000011 has no access to StateStore null as the store is not connected to the processor.

If you add stores manually via '.addStateStore()'
make sure to connect the added store to the processor by providing the processor name to '.addStateStore()'
or connect them via '.connectProcessorAndStateStores()'.

DSL users need to provide the store name to '.process()', '.transform()', or '.transformValues()'
to connect the store to the corresponding operator,
or they can provide a StoreBuilder by implementing the stores() method on the Supplier itself.
If you do not add stores manually, please file a bug report at https://issues.apache.org/jira/projects/KAFKA.

代码A

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.*;
import org.apache.kafka.streams.processor.AbstractProcessor;
import org.apache.kafka.streams.processor.ProcessorContext;
import org.apache.kafka.streams.processor.PunctuationType;
import org.apache.kafka.streams.state.KeyValueIterator;
import org.apache.kafka.streams.state.KeyValueStore;
import org.apache.kafka.streams.state.Stores;

import java.time.Duration;
import java.util.Properties;

public class StreamProcessingApp {

    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "stream-processing-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());

        StreamsBuilder builder = new StreamsBuilder();
        createTopology(builder);

        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        streams.start();

        Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
    }

    public static void createTopology(StreamsBuilder builder) {
        KStream<String, String> stream1 = builder.stream("input-topic-1");
        KStream<String, String> stream2 = builder.stream("input-topic-2");

        // Define and register the state store
        final String stateStoreName = "join-store";
        builder.addStateStore(
            Stores.keyValueStoreBuilder(
                Stores.persistentKeyValueStore(stateStoreName),
                Serdes.String(),
                Serdes.String()
            )
        );

        KStream<String, String> joinedStream = stream1.outerJoin(
                stream2,
                (value1, value2) -> {
                    if (value1 == null) {
                        return "null-" + value2;
                    } else if (value2 == null) {
                        return value1 + "-null";
                    }
                    return value1 + "-" + value2;
                },
                JoinWindows.of(Duration.ofMinutes(5)),
                StreamJoined.with(Serdes.String(), Serdes.String(), Serdes.String())
        );

        joinedStream.to("output-topic");

        // Process the joined stream and utilize the state store
        joinedStream.process(() -> new AbstractProcessor<String, String>() {
            private KeyValueStore<String, String> stateStore;

            @Override
            public void init(ProcessorContext context) {
                super.init(context);
                this.stateStore = (KeyValueStore<String, String>) context.getStateStore(stateStoreName);
                context.schedule(Duration.ofMinutes(1), PunctuationType.WALL_CLOCK_TIME, timestamp -> {
                    // Check for values that did not join
                    KeyValueIterator<String, String> iterator = this.stateStore.all();
                    while (iterator.hasNext()) {
                        var entry = iterator.next();
                        if (entry.value.endsWith("-null") || entry.value.startsWith("null-")) {
                            System.out.println("Unmatched record found: Key=" + entry.key + ", Value=" + entry.value);
                            // Custom handling logic for unmatched records
                        }
                    }
                    iterator.close();
                });
            }

            @Override
            public void process(String key, String value) {
                // Store each joined record in the state store
                this.stateStore.put(key, value);
            }

            @Override
            public void close() {
                // Cleanup
            }
        }, stateStoreName);  // Connect the state store to the processor
    }
}

2024.08.02 下午3:22 更新

已针对该问题做初步调研,大部分方案要求在添加Processor时传入StateStore名称(stream.process(..., "storeName")),但已按此操作仍出现相同错误:

Caused by: org.apache.kafka.streams.errors.StreamsException: Processor KSTREAM-TRANSFORM-0000000002 has no access to StateStore my-store as the store is not connected to the processor...

调研参考内容:

  • Spring Cloud Stream Binder Kafka 3.x 无法将自定义存储连接到Transformer
  • Spring Kafka - 添加的存储无法从流处理中访问
  • 为什么我的Kafka Transformer的StateStore不可访问?

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 02:55:11