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

Apple M1环境下KStream关联KTable触发RocksDB列族不存在错误

问题:Kafka Streams KStream与KTable LeftJoin报错(Apple M1环境)

我在对两个流执行leftJoin操作时,最初用两个KStream关联正常,但把第二个流转为KTable后出现错误。相关代码如下:

@Bean
public KafkaStreams kafkaStreams() throws IOException {
        final Properties props = configureKafkaStreamsProperties();
            
        ObjectMapper mapper = new ObjectMapper();
        mapper.disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES);

        final StreamsBuilder builder = new StreamsBuilder();

        // 1st Structured stream
        KStream<String, String> firstStream = builder.stream("topic-1", Consumed.with(Serdes.String(), Serdes.String()));


        KStream<String, String> firstStreamTransformed = firstStream.map((k, v) -> {
                        try {
                                InputModelOne model = mapper.readValue(v, InputModelOne.class);
                                return new KeyValue<>(model.getId(), v);
                        } catch (Exception e) {
                                logger.error(e.toString());
                                return new KeyValue<>(k, v);
                        }

         });


        // Second stream
        KStream<String, String> secondStream = builder.stream("topic-2",
                                Consumed.with(Serdes.String(), Serdes.String()));

        KStream<String, String> secondStreamTransformed = secondStream.map((k, v) -> {
                        try {
                                InputModelTwo model = mapper.readValue(v, InputModelTwo.class);
                                return new KeyValue<>(model.getId(), v);

                        } catch (Exception e) {
                                logger.error(e.toString());
                                return new KeyValue<>(k, v);
                        }
        });

        // Build KTable from second topic
        KTable<String, String> secondTable = secondStreamTransformed.toTable(Materialized.as("topic-2-table"));


        // Valuejoiner
        ValueJoiner<String, String, String> joiner = (one, two) -> {

                try {
                                
                    InputModelOne modelOne = mapper.readValue(one, InputModelOne.class);
                    InputModelTwo modelTwo = new InputModelTwo();

                    // Create output object with properties
                    OutputModel out = new OutputModel(modelOne.getId());
                    out.setOneTimestamp(modelOne.getTimestamp());
                    out.setTwoTimestamp(modelTwo.getTimestamp());

                    return mapper.writeValueAsString(out);
                    } catch (JsonProcessingException e) {
                                // TODO Auto-generated catch block
                                e.printStackTrace();
                                return null;
                    }
         };

         KStream<String, String> joined = firstStreamTransformed.leftJoin(secondTable,
                                joiner);


         joined.to("joined-topics", Produced.with(Serdes.String(), Serdes.String()));
}

报错信息如下:

org.apache.kafka.streams.errors.ProcessorStateException: Error opening store joined-topics at location /var/folders/lx/dz_x9j5d7lz4mfymgzkcn7wr0000gn/T/kafka-streams/streams-pipe/2_0/rocksdb/joined-topics
        at org.apache.kafka.streams.state.internals.RocksDBTimestampedStore.openRocksDB(RocksDBTimestampedStore.java:87) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.state.internals.RocksDBStore.openDB(RocksDBStore.java:186) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.state.internals.RocksDBStore.init(RocksDBStore.java:254) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.state.internals.WrappedStateStore.init(WrappedStateStore.java:55) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.state.internals.ChangeLoggingKeyValueBytesStore.init(ChangeLoggingKeyValueBytesStore.java:55) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.state.internals.WrappedStateStore.init(WrappedStateStore.java:55) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.state.internals.CachingKeyValueStore.init(CachingKeyValueStore.java:75) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.state.internals.WrappedStateStore.init(WrappedStateStore.java:55) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.state.internals.MeteredKeyValueStore.lambda$init$1(MeteredKeyValueStore.java:122) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:884) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.state.internals.MeteredKeyValueStore.init(MeteredKeyValueStore.java:122) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.processor.internals.ProcessorStateManager.registerStateStores(ProcessorStateManager.java:201) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.processor.internals.StateManagerUtil.registerStateStores(StateManagerUtil.java:103) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.processor.internals.StreamTask.initializeIfNeeded(StreamTask.java:216) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.processor.internals.TaskManager.tryToCompleteRestoration(TaskManager.java:433) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.processor.internals.StreamThread.initializeAndRestorePhase(StreamThread.java:849) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:731) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:583) ~[kafka-streams-2.8.0.jar:na]
        at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:556) ~[kafka-streams-2.8.0.jar:na]
Caused by: org.rocksdb.RocksDBException: Column family not found: keyValueWithTimestamp
        at org.rocksdb.RocksDB.open(Native Method) ~[rocksdbjni-6.29.4.1.jar:na]
        at org.rocksdb.RocksDB.open(RocksDB.java:306) ~[rocksdbjni-6.29.4.1.jar:na]
        at org.apache.kafka.streams.state.internals.RocksDBTimestampedStore.openRocksDB(RocksDBTimestampedStore.java:75) ~[kafka-streams-2.8.0.jar:na]
        ... 18 common frames omitted

我使用Docker在本地运行Kafka和ZooKeeper,设备为Apple M1,希望得到解决建议以继续在Mac上开发。


解决建议
  • 清理本地状态存储:报错核心是RocksDB找不到指定列族,这通常是因为之前运行的拓扑生成的旧状态数据与当前拓扑不兼容。找到报错中的状态存储目录/var/folders/lx/dz_x9j5d7lz4mfymgzkcn7wr0000gn/T/kafka-streams/,删除整个目录后重启应用。

  • 显式指定KTable状态存储的Serdes:创建KTable时,显式声明key和value的序列化器,避免自动推断导致的不匹配问题。修改代码如下:

KTable<String, String> secondTable = secondStreamTransformed.toTable(
    Materialized.<String, String>as("topic-2-table")
        .withKeySerde(Serdes.String())
        .withValueSerde(Serdes.String())
);
  • 升级Kafka Streams版本:你当前使用的kafka-streams-2.8.0在Apple M1的ARM架构上,RocksDB JNI可能存在适配问题。建议升级到3.0及以上的稳定版本,新版本对ARM架构支持更完善。

  • 使用适配ARM架构的Docker Kafka镜像:确保本地Docker运行的Kafka是支持arm64架构的镜像,比如confluentinc/cp-kafka:7.4.0,避免架构不兼容引发的潜在异常。

  • 处理LeftJoin中的空值场景:LeftJoin中KTable的对应值可能为null,当前代码直接new InputModelTwo(),建议先判断非空再反序列化,避免后续逻辑出现问题:

ValueJoiner<String, String, String> joiner = (one, two) -> {
    try {
        InputModelOne modelOne = mapper.readValue(one, InputModelOne.class);
        InputModelTwo modelTwo = null;
        if (two != null) {
            modelTwo = mapper.readValue(two, InputModelTwo.class);
        } else {
            modelTwo = new InputModelTwo();
        }
        OutputModel out = new OutputModel(modelOne.getId());
        out.setOneTimestamp(modelOne.getTimestamp());
        out.setTwoTimestamp(modelTwo != null ? modelTwo.getTimestamp() : null);
        return mapper.writeValueAsString(out);
    } catch (JsonProcessingException e) {
        e.printStackTrace();
        return null;
    }
};

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 23:50:43