【问题标题】:Kafka Streams Binding: Cannot get state store because the stream thread is PARTITIONS_ASSIGNED, not RUNNINGKafka Streams Binding:无法获取状态存储,因为流线程是 PARTITIONS_ASSIGNED,而不是 RUNNING
【发布时间】: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


【解决方案1】:

您需要启用运行状况指示器,然后编写侦听器以在第二次 RUNNING 之前监视它的 REBALANCE。

参考这个 Health Indicator

【讨论】:

  • 我知道健康指标。你能否提供更多关于编写监听器来监控的细节?可能是一些有用的参考?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-04-07
  • 2019-07-03
  • 1970-01-01
  • 2018-11-14
  • 2018-10-24
  • 2021-01-04
相关资源
最近更新 更多