Kafka Streams绑定:Store(prod id count Store)对Streams实例不可用



我正在尝试使用Spring Cloud Stream Kafka Binder在Kafka Streams上运行交互式查询。我坚持从InteractiveQueryService检索键值存储。我总是在我的代码上,甚至在示例中的代码上得到相同的错误:

https://github.com/spring-cloud/spring-cloud-stream-samples/tree/main/kafka-streams-samples/kafka-streams-interactive-query-basic

java.lang.IllegalStateException: Error retrieving state store: prod-id-count-store
at org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService.lambda$getQueryableStore$0(InteractiveQueryService.java:115) ~[spring-cloud-stream-binder-kafka-streams-3.2.5.jar:3.2.5]
at org.springframework.retry.support.RetryTemplate.doExecute(RetryTemplate.java:329) ~[spring-retry-1.3.1.jar:na]
at org.springframework.retry.support.RetryTemplate.execute(RetryTemplate.java:209) ~[spring-retry-1.3.1.jar:na]
at org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService.getQueryableStore(InteractiveQueryService.java:88) ~[spring-cloud-stream-binder-kafka-streams-3.2.5.jar:3.2.5]
at com.nowatel.teacher.service.KafkaStreamsInteractiveQueryApplication$InteractiveProductCountApplication.printProductCounts(KafkaStreamsInteractiveQueryApplication.java:71) ~[main/:na]
at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) ~[na:na]
at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77) ~[na:na]
at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[na:na]
at java.base/java.lang.reflect.Method.invoke(Method.java:568) ~[na:na]
at org.springframework.scheduling.support.ScheduledMethodRunnable.run(ScheduledMethodRunnable.java:84) ~[spring-context-5.3.23.jar:5.3.23]
at org.springframework.scheduling.support.DelegatingErrorHandlingRunnable.run(DelegatingErrorHandlingRunnable.java:54) ~[spring-context-5.3.23.jar:5.3.23]
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) ~[na:na]
at java.base/java.util.concurrent.FutureTask.runAndReset$$$capture(FutureTask.java:305) ~[na:na]
at java.base/java.util.concurrent.FutureTask.runAndReset(FutureTask.java) ~[na:na]
at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:305) ~[na:na]
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) ~[na:na]
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) ~[na:na]
at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na]
Caused by: org.apache.kafka.streams.errors.UnknownStateStoreException: Store (prod-id-count-store) not available to Streams instance
at org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService.getStateStoreFromKafkaStreams(InteractiveQueryService.java:126) ~[spring-cloud-stream-binder-kafka-streams-3.2.5.jar:3.2.5]
at org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService.lambda$getQueryableStore$0(InteractiveQueryService.java:109) ~[spring-cloud-stream-binder-kafka-streams-3.2.5.jar:3.2.5]
... 17 common frames omitted

代码来自示例类:

https://github.com/spring-cloud/spring-cloud-stream-samples/blob/main/kafka-streams-samples/kafka-streams-interactive-query-basic/src/main/java/kafka/streams/product/tracker/KafkaStreamsInteractiveQueryApplication.java

我的application.properties文件:

spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000
spring.cloud.stream.bindings.process-out-0.destination=product-counts
spring.cloud.stream.bindings.process-in-0.destination=products
spring.application.name=kafka-streams-iq-basic-sample

唯一的区别是Spring Boot版本、Spring Cloud版本和Java 17。

我的build.gradle文件:

plugins {
id 'org.springframework.boot' version "2.7.4"
id 'io.spring.dependency-management' version "1.0.14.RELEASE"
id 'java'
}
group = 'com.example'
version = '0.0.1-SNAPSHOT'
sourceCompatibility = javaVersion
repositories {
mavenCentral()
}
dependencies {
implementation 'org.springframework.boot:spring-boot-starter-webflux'
implementation 'org.apache.kafka:kafka-streams'
implementation 'org.springframework.cloud:spring-cloud-stream'
implementation 'org.springframework.cloud:spring-cloud-stream-binder-kafka'
implementation 'org.springframework.cloud:spring-cloud-stream-binder-kafka-streams'
implementation 'org.springframework.kafka:spring-kafka'
testImplementation 'org.springframework.boot:spring-boot-starter-test'
testImplementation 'org.springframework.kafka:spring-kafka-test'
testImplementation 'io.projectreactor:reactor-test'
}
dependencyManagement {
imports {
mavenBom "org.springframework.cloud:spring-cloud-dependencies:2021.0.4"
}
}
tasks.named('test') {
useJUnitPlatform()
}

Kafka版本:Kafka_2.13-3.2.3

你能告诉我做错了什么吗,因为我浪费了两天时间,没有任何想法。

修复程序在main分支中,但尚未在3.2.x中。但是,您应该能够通过设置来解决此问题:

spring.cloud.stream.kafka.streams.binder.configuration.application.server=<server>:<port>

我认为这与这个问题有关:https://github.com/spring-cloud/spring-cloud-stream/issues/2523

你能对照活页夹的最新3.2.x快照试试吗?(3.2.6-SNAPSHOT(。我们很快就会发布。

您只需要更新绑定器依赖项-spring-cloud-stream-binder-kafka-streams(到3.2.6-SNAPSHOT(。

最新更新