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类构建查询参数,以此提供更灵活的查询配置能力。
修改步骤
- 导入
StoreQueryParameters类(若尚未导入):
import org.apache.kafka.streams.StoreQueryParameters;
- 修改
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
相关产品推荐
相关产品推荐

