【问题标题】:How to add a StateStore using StateStoreBuilder in a Spring Cloud Stream Kafka Streams application如何在 Spring Cloud Stream Kafka Streams 应用程序中使用 StateStoreBuilder 添加 StateStore
【发布时间】: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


    【解决方案1】:

    由于我的 Scs 版本 (Fishtown SR3),我没有看到文档化的过程,但好消息是,从 Germantown 开始就可以以声明方式创建 State Store:

    const val DEDUP_STORE = "dedup-store"
    
    @EnableBinding(KafkaStreamsProcessor::class)
    class FraudKafkaStreamsConfiguration {
    
        @KafkaStreamsStateStore(name = DEDUP_STORE, type = KafkaStreamsStateStoreProperties.StoreType.KEYVALUE)
        @StreamListener
        @SendTo("output")
        fun process(@Input("input") input: KStream<String, TransferEmitted>): KStream<String, TransferEmitted> {
            return input.transform(TransformerSupplier { DeduplicationTransformer() }, DEDUP_STORE)
    
        }
    
    }
    

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-12-09
    • 2023-03-10
    • 2020-07-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多