【问题标题】:Kafka stream enrichment when we must fetch data on the fly当我们必须动态获取数据时,Kafka 流丰富
【发布时间】:2018-02-07 16:17:25
【问题描述】:

我们有一个传入事件流:

{ user: 123, itemid: 'abc', ... }
{ user: 456, itemid: 'abc', ... }
{ user: 123, itemid: 'def', ... }
{ user: 789, itemid: 'xyz', ... }
{ user: 123, itemid: 'xyz', ... }

等等。我们需要使用通过项目 ID(Web 服务调用)查找的数据来丰富这些,但是 查找可能很慢。另一方面,项目数据是相当静态的, 所以缓存一天不会有问题。

现在,标准解决方案是将项目数据发布到 KTable 并加入,但是... 我们无法获取所有项目 - 我们从后端系统获取这些项目,不能只是“获取所有”。 并且不会有更新通知。想象一下,项目 ID 是书的 ISBN 号 - 我们不能要求世界上所有书籍的清单。而图书信息只会慢慢改变 (例如,有人修复了分类错误),我们不需要为此类更改重新发布事件。

那么,我们该怎么做呢? 我目前的想法是我们在 itemid 上进行分区,与项目 KTable 进行左连接,然后在值连接器中, 如果正确的值为 null,则进行 Web 服务查找,返回结果并将结果发布到项目 KTable 主题。 但这会阻止我们查找 'abc' 两次吗?

就规模而言,我们当前的数据集包含大约 150 万个单独的项目。

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    基本想法似乎是合理的,但我不会使用KTable,而是使用维护状态的自定义有状态Transformer()transform() 只有一个输入流,您可以查找 Transformers 状态。如果查找失败,您将转到 RCP,将其放入存储中并为您的输出记录执行连接。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2012-01-22
      • 2020-11-20
      • 2022-06-12
      • 2019-05-31
      • 1970-01-01
      • 1970-01-01
      • 2011-02-11
      相关资源
      最近更新 更多