【问题标题】:how to make a Spring cloud stream Processor Function that takes kstream (Kafka stream) to work timely basis?如何制作一个以kstream(Kafka流)及时工作的Spring云流处理器功能?
【发布时间】:2021-07-13 23:42:11
【问题描述】:

所以我的 云流处理器功能 从一个 Kafka 主题 1 读取消息并将消息生成到另一个 Kafka 主题 2 的场景。但是这个过程必须及时运行,比如函数应该等待 5 分钟,然后它应该启动(消费 n 生产)1 分钟,然后在 1 分钟后再次等待 5 分钟。谁能帮我看看怎么做?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams spring-cloud-stream spring-cloud-stream-binder-kafka


    【解决方案1】:

    Spring Cloud Stream 函数是事件驱动的,因此您无法控制消费时间。一旦它可供活页夹使用,它就会被消耗掉。如果这是一个硬性要求,我建议从 Spring 编写一个带有 Scheduled 注释的常规 bean,并使用 Spring for Apache Kafka 从第一个主题消费,然后立即生产到第二个主题。您可以使用来自 Sping Kafka 的 KafkaListener 和 KafkaTemplate(或来自 Spring Cloud Stream 的 StreamBridge 用于发送到 Kafka)。一旦记录在第二个主题中,您仍然可以使用 Kafka Streams binder 进行进一步的流处理。

    【讨论】:

    • 这个问题说的是不同的主题,没有关于不同的集群
    • @sobychacko - 感谢您的回复。这里我们有一个集群。它涉及两个不同的主题。我必须从一个主题消费并将其生产到另一个主题。但这整个过程我需要及时控制。
    • 抱歉没有正确阅读问题。我更新了上面的答案。
    猜你喜欢
    • 2021-07-02
    • 1970-01-01
    • 2021-07-28
    • 2015-09-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-11-29
    • 1970-01-01
    相关资源
    最近更新 更多