【问题标题】:KStreams - KTable startup phaseKStreams - KTable 启动阶段
【发布时间】:2018-03-18 21:09:35
【问题描述】:

KStreams - KTable 连接的工作方式非常简单:每次在流上发出新样本时,都会在表上执行键查找。

这会在瞬态阶段产生意外行为吗?我们有这样的拓扑:

  1. 一个KStream A,我们执行 selectKey 将其转换为 Stream A1
  2. 一个KStream B,我们groupBy然后减少,变成一个KTable B1

在启动时,我们在 A 上发布两条记录,在 B 上发布两条记录,这样在 A 上的 selectKey 和 B 上的 groupBy + reduce 之后,key 就会匹配。但是,我们注意到有时 A1 和 B1 之间的内部连接的样本会失败,而我们会丢失一些我们期望的输出。

确保没有更新丢失的正确拓扑是什么?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    KStream-KTable 加入同步是最好的努力。我们致力于改进,为 1.2 版本提供更好的保证。 Atm,你无能为力。

    如果您需要严格的保证,则需要使用transform() 而不是join() 实现自己的流表连接运算符。您可以将 KTable 存储连接到 Transformer,并为连接查找放置自定义逻辑。

    【讨论】:

    • 我在变压器里面做什么?
    • 我的流表连接是键值查找。因此,对于 tranform() 的每个输入记录,您在 KTable 存储中执行 get()。或者,当然,您需要为 KTable 尚未更新以便稍后重试的情况设置一些逻辑。
    • 你将如何在流处理的上下文中重构这个重试逻辑?
    • 这完全取决于您...这取决于我要说的用例。 (这就是为什么 KStream-KTable 连接首先是 atm 的最大努力:“kstream 处理器”会进行键值查找,但如果没有匹配则不会重试——即,如果 KTable 更新延迟,你会错过更新)。
    • 您好@MatthiasJ.Sax,感谢您的回复!这在 2.3 中是否仍然相同?更具体地说是在 ktable 和 ktable 内连接之间?我看到了同样的行为,在我采用变换方法之前,我想知道我是否走错了方向,但体验到的效果是一样的!
    猜你喜欢
    • 2012-08-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-04-21
    相关资源
    最近更新 更多