gpt4 book ai didi

apache-kafka - 为什么我在检索商店进行查询时偶尔会收到 InvalidStateStoreException PARTITIONS_REVOKED,而不是 RUNNING?

转载 作者:行者123 更新时间:2023-12-04 15:55:17 25 4
gpt4 key购买 nike

我正在访问一个状态存储来查询它并且不得不用一个 try/catch block 包装 store() 语句来重试它,因为有时我会得到这个异常:

org.apache.kafka.streams.errors.InvalidStateStoreException: Cannot get state store customers-store because the stream thread is PARTITIONS_REVOKED, not RUNNING
at org.apache.kafka.streams.state.internals.StreamThreadStateStoreProvider.stores(StreamThreadStateStoreProvider.java:49)
at org.apache.kafka.streams.state.internals.QueryableStoreProvider.getStore(QueryableStoreProvider.java:57)
at org.apache.kafka.streams.KafkaStreams.store(KafkaStreams.java:1053)
at com.codependent.kafkastreams.customer.service.CustomerService.getCustomer(CustomerService.kt:75)
at com.codependent.kafkastreams.customer.service.CustomerServiceKt.main(CustomerService.kt:108)

这是用于检索商店的代码(完整代码在 github repo 上):

fun getCustomer(id: String): Customer? {
var keyValueStore: ReadOnlyKeyValueStore<String, Customer>? = null
while(keyValueStore == null) {
try {
keyValueStore = streams.store(CUSTOMERS_STORE, QueryableStoreTypes.keyValueStore<String, Customer>())
} catch (ex: InvalidStateStoreException) {
ex.printStackTrace()
}
}
val customer = keyValueStore.get(id)
return customer
}

这是主程序:

fun main(args: Array<String>) {
val customerService = CustomerService("main", "localhost:9092")
customerService.initializeStreams()
customerService.createCustomer(Customer("53", "Joey"))
val customer = customerService.getCustomer("53")
println(customer)
customerService.stopStreams()
}

在先前的执行完成后,异常发生随机运行程序几次。注意:我没有对正在执行的 Kafka 集群做任何事情并使用它的默认配置。

最佳答案

在您访问商店时,Kafka Streams 应用程序正在重新平衡,此时无法访问状态商店。您希望确保仅在应用程序状态为 RUNNING 而不是 REBALANCING 时查询商店。

您可以做的是在尝试像这样从商店读取之前检查应用程序的状态:

if(streams.state() == State.RUNNING) {
keyValueStore = streams.store(...);
val customer = keyValueStore.get(id);
return customer;
}

还有一个 KafkaStreams.setStateListener 方法可以用来注册一个 KafkStreams.StateListener 实现。每次应用程序更改其状态时,都会调用 StateListener.onChange 方法。

关于apache-kafka - 为什么我在检索商店进行查询时偶尔会收到 InvalidStateStoreException PARTITIONS_REVOKED,而不是 RUNNING?,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/52016385/

25 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com