【问题标题】:Kafka Streams: Punctuate vs ProcessKafka Streams:标点符号与流程
【发布时间】:2018-06-09 17:48:24
【问题描述】:

在流应用程序中的单个任务中,以下两个方法是否独立运行(意味着当方法“process”处理来自上游源的传入消息时,方法“punctuate”也可以基于指定 schedule 和 WALL_CLOCK_TIME 作为 PunctuationType?)或者它们是否共享同一个线程,因此它要么是在给定时间运行的线程,如果是这样,如果 process 方法不断从上游源获取消息,那么 punctuate 方法永远不会被调用?

  • Processor.process(K 键,V 值)
    使用给定的键和值处理记录。

  • ProcessorContext.schedule(长间隔,PunctuationType 类型,Punctuator 回调)
    为处理器安排定期操作。

另外,请澄清在标点法中分区 id 值为 -1 是什么意思。标点方法不是特定于任何分区吗?

  • int ProcessorContext.partition()
    返回当前输入记录的分区id;如果它不可用(例如,如果从标点调用调用此方法),则可能为 -1

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    这两种方法都在一个线程中执行。如果有输入数据,基于挂钟的punctuate() 将被独立调用:在调用process() 之间,线程检查系统时间并在必要时调用punctuate()

    对于分区信息:是的,标点符号与分区无关。当然,标点符号是特定于任务的,但是,一个任务可能有多个输入分区(例如,如果它执行mergejoin),因此不清楚要传递哪些分区信息。为简单起见,单个分区大小写的处理方式与多分区大小写相同,标点符号与分区分离。

    【讨论】:

    • 嗨,马修,谢谢。假设流应用程序没有合并或连接,并且在具有 6 个分区的主题中,如果调用 punctuate() 并且如果我打印 context.TaskId(),它会反映单个任务 0(左)和右侧的相应分区(0_1;0_2;0_3;0_4;0_5;0_6)。使用TaskId方法判断标点中任务对应哪个分区有效吗?
    • 在此设置和当前实现中,任务 ID 的第二个数字与分区号相同。但是,此命名架构不是公共合同的一部分,并且可能会在未来版本中更改,恕不另行通知。因此,不建议编写依赖于此的代码,因为如果您升级到较新的版本,它可能会崩溃!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-06-14
    • 2020-07-30
    • 2020-01-20
    • 1970-01-01
    相关资源
    最近更新 更多