【发布时间】:2021-12-29 09:16:58
【问题描述】:
我使用 Spring Cloud Stream Kafka Binding 的基础架构进行流处理。
这里是主要的 KafkaStream 方法:
private static final String READ_STORE_NAME = "readMessageStore";
@Bean
public Consumer<KStream<String, ContentModel>> filter() {
return in -> in.groupByKey(Grouped.with(Serdes.String(), new JsonSerde<>(ContentModel.class)))
.reduce(this::mergeObjects, Materialized.as(READ_STORE_NAME));
}
private ContentModel mergeObjects(final ContentModel reducer, final ContentModel materialized) {
return reducer;
}
绑定配置:
spring.cloud.stream.bindings.filter-in-0:
destination: topic
我部署了2个应用实例,用ApplicationId分隔如下:
spring.cloud.stream.kafka.streams.binder.functions.filter.applicationId = <Id>
但在部署阶段我收到:
IllegalStateException:检索状态存储时出错:j readMessageStore
原因:org.apache.kafka.streams.errors.InvalidStateStoreException:无法获取状态存储 readMessageStore,因为流线程是 PARTITIONS_ASSIGNED,而不是 RUNNING
从Unable to open store for Kafka streams because invalid state 我发现,我需要(显然)在查询之前等待第二个 RUNNING 状态。
但是考虑到这种情况下的所有状态转换都是“在后台”实现的,我该如何在 Spring Cloud Stream Kafka Binding 环境中设置此延迟?
【问题讨论】:
-
您使用的是哪个代理?从您链接的另一个 SO 线程来看,如果您使用的是经纪人 >= 2.2.0,那么这已经修复了吗?你试过了吗?
-
在 kafka 流绑定器的情况下,据我所知,经纪人是“幕后”。
标签: apache-kafka apache-kafka-streams spring-cloud-stream