【问题标题】:Setting Window[Hopping, Tumbling..etc] Parallelism in Kafka Streams设置窗口[Hopping, Tumbling..etc] Kafka Streams 中的并行度
【发布时间】:2017-11-29 14:26:34
【问题描述】:

我正在编写 kafka Streams 中的 Hopping Window Code,其中 minMaxCalculator() 在流按键分组后计算流内的最小值和最大值。

 KTable<Windowed<String>, aggrTest> WinMinMax = Records.groupByKey().aggregate(new aggrTestInitilizer(), 
        new minMaxCalculator()
        , TimeWindows.of(TimeUnit.SECONDS.toMillis(5)).advanceBy(TimeUnit.SECONDS.toMillis(1)),aggrMessageSerde,"aggr-test");

一旦我按键分组,我想并行处理为所有键生成的窗口,即使有一个 kafka 分区。我们应该怎么做?我在哪里可以设置与窗口对应的并行度

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    并行性基于输入分区,不能与它们不同。因此,您无法设置任何参数。

    但是,您可以创建具有所需分区数的主题,并使用 through() 将其用于手动重新分区:

    stream.through("multi-partition-topic").groupByKey()...
    

    查看文档了解更多详情:

    【讨论】:

    • 小说明:输入分区的数量决定了最大的并行度。您当然可以运行比此最大值更少的应用程序实例(或线程)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-11-10
    • 2020-09-05
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多