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

Kafka Streams从2.x升级到3.3.2后store方法调用报错

Kafka Streams 3.3.2版本store方法调用不兼容问题修复

问题背景

将Kafka依赖从2.x升级到3.3.2后,调用KafkaStreams.store()方法出现编译错误,无法匹配原有的方法参数。

错误信息

store(org.apache.kafka.streams.StoreQueryParameters<T>)' in 'org.apache.kafka.streams.KafkaStreams' cannot be applied to '(java.lang.String, org.apache.kafka.streams.state.QueryableStoreType<org.apache.kafka.streams.state.ReadOnlyKeyValueStore<java.lang.Object,java.lang.Object>>)'

相关代码

QueryServer.java

import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.state.HostInfo;
import org.apache.kafka.streams.state.QueryableStoreTypes;
import org.apache.kafka.streams.state.ReadOnlyKeyValueStore;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import spark.Spark;

import javax.ws.rs.client.Client;
import javax.ws.rs.client.ClientBuilder;
import java.util.ArrayList;
import java.util.List;

class QueryServer {
    private static final Logger logger = LogManager.getLogger();
    private final String NO_RESULTS = "No Results Found";
    private final String APPLICATION_NOT_ACTIVE = "Application is not active. Try later.";
    private final KafkaStreams streams;
    private Boolean isActive = false;
    private final HostInfo hostInfo;
    private Client client;

    QueryServer(KafkaStreams streams, String hostname, int port) {
        this.streams = streams;
        this.hostInfo = new HostInfo(hostname, port);
        client = ClientBuilder.newClient();
    }

    void setActive(Boolean state) {
        isActive = state;
    }

    private List<KeyValue<String, String>> readAllFromLocal() {

        List<KeyValue<String, String>> localResults = new ArrayList<>();
        ReadOnlyKeyValueStore<String, String> stateStore =
            streams.store(AppConfigs.stateStoreName, QueryableStoreTypes.keyValueStore());

        stateStore.all().forEachRemaining(localResults::add);
        return localResults;
    }

    void start() {
        logger.info("Starting Query Server at http://" + hostInfo.host() + ":" + hostInfo.port()
            + "/" + AppConfigs.stateStoreName + "/all");

        Spark.port(hostInfo.port());

        Spark.get("/" + AppConfigs.stateStoreName + "/all", (req, res) -> {

            List<KeyValue<String, String>> allResults;
            String results;

            if (!isActive) {
                results = APPLICATION_NOT_ACTIVE;
            } else {
                allResults = readAllFromLocal();
                results = (allResults.size() == 0) ? NO_RESULTS
                    : allResults.toString();
            }
            return results;
        });

    }

    void stop() {
        client.close();
        Spark.stop();
    }
}

MainApp.java

public class StreamingTableApp {
    private static final Logger logger = LogManager.getLogger();

    public static void main(final String[] args) {

        final Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, AppConfigs.applicationID);
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, AppConfigs.bootstrapServers);
        props.put(StreamsConfig.STATE_DIR_CONFIG, AppConfigs.stateStoreLocation);
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder streamsBuilder = new StreamsBuilder();
        KTable<String, String> KT0 = streamsBuilder.table(AppConfigs.topicName);
        KT0.toStream().print(Printed.<String, String>toSysOut().withLabel("KT0"));

        KTable<String, String> KT1 = KT0.filter((k, v) -> k.matches(AppConfigs.regExSymbol) && !v.isEmpty(),
            Materialized.as(AppConfigs.stateStoreName));
        KT1.toStream().print(Printed.<String, String>toSysOut().withLabel("KT1"));

        KafkaStreams streams = new KafkaStreams(streamsBuilder.build(), props);

        //Query Server
        QueryServer queryServer = new QueryServer(streams, AppConfigs.queryServerHost, AppConfigs.queryServerPort);
        streams.setStateListener((newState, oldState) -> {
            logger.info("State Changing to " + newState + " from " + oldState);
            queryServer.setActive(newState == KafkaStreams.State.RUNNING && oldState == KafkaStreams.State.REBALANCING);
        });

        streams.start();
        queryServer.start();

        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            logger.info("Shutting down servers");
            queryServer.stop();
            streams.close();
        }));

    }
}

解决方案

原因分析

Kafka Streams 3.0及以上版本中,原有的store(String storeName, QueryableStoreType<T> storeType)方法已被废弃,官方改为使用StoreQueryParameters类构建查询参数,以此提供更灵活的查询配置能力。

修改步骤

  1. 导入StoreQueryParameters类(若尚未导入):
import org.apache.kafka.streams.StoreQueryParameters;
  1. 修改QueryServer.java中readAllFromLocal()方法的store获取逻辑:
private List<KeyValue<String, String>> readAllFromLocal() {

    List<KeyValue<String, String>> localResults = new ArrayList<>();
    // 构建查询参数
    StoreQueryParameters<ReadOnlyKeyValueStore<String, String>> queryParams = 
        StoreQueryParameters.fromNameAndType(
            AppConfigs.stateStoreName, 
            QueryableStoreTypes.keyValueStore()
        );
    // 使用新方式获取state store
    ReadOnlyKeyValueStore<String, String> stateStore = streams.store(queryParams);

    stateStore.all().forEachRemaining(localResults::add);
    return localResults;
}

说明

StoreQueryParameters.fromNameAndType()方法会根据传入的存储名称和查询类型生成对应的查询参数对象,完全替代旧版本的双参数store()方法调用,同时支持后续扩展更多查询配置(如指定查询副本等)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 13:15:17