【发布时间】:2019-10-30 23:28:52
【问题描述】:
原生 Kafka API 允许我们create and add a state store using the StreamsBuilder:
final StreamsBuilder builder = new StreamsBuilder();
...
final StoreBuilder<WindowStore<String, Long>> dedupStoreBuilder = Stores.windowStoreBuilder(
Stores.persistentWindowStore(storeName,
retentionPeriod,
windowSize,
false
),
Serdes.String(),
Serdes.Long());
builder.addStateStore(dedupStoreBuilder);
我想使用 Spring Cloud Streams 做同样的事情,但不知道如何访问 StreamsBuilder 以添加商店。
我已尝试按照doc 中的说明检索StreamsBuilderFactoryBean,希望可以从中获取StreamsBuilder 对象,但该bean 似乎不可用:
@EnableBinding(KafkaStreamsProcessor::class)
class FraudKafkaStreamsConfiguration(private val context: ApplicationContext) {
@StreamListener
@SendTo("output")
fun process(@Input("input") input: KStream<String, TransferEmitted>): KStream<String, TransferEmitted> {
val streamsBuilderFactoryBean = context.getBean("&stream-builder-process", StreamsBuilderFactoryBean::class.java)
...
return xxx
}
}
原因: org.springframework.beans.factory.NoSuchBeanDefinitionException: 否 名为“stream-builder-process”的 bean 可用
无论如何,我什至不确定这是正确的做法。那么,我们如何以编程方式创建StateStore?
【问题讨论】:
标签: apache-kafka spring-cloud apache-kafka-streams spring-cloud-stream