【问题标题】:Aggregate messages using multiple fields in the message value使用消息值中的多个字段聚合消息
【发布时间】:2019-05-27 16:05:53
【问题描述】:

我有一个 Kafka 主题,其中包含多个不同用户的多个用户信息事件。 我试图弄清楚如何使用值中的多个字段将这些聚合在一起。

例如:

输入主题:

1:{"SSN":"123456"}
2:{"twitterHandle":"elvis"}
3:{"SSN":"123456","twitterHandle":"elvis","accountNum": "111111"}
4:{"SSN":"123456"}
5:{"SSN":"000000"}
6:{"twitterHandle":"foo"}
7:{"SSN":"000000","twitterHandle":"foo"}
8:{"SSN":"000000"}

我想要一个输出主题(聚合):

{"SSN":"123456","twitterHandle":"elvis","accountNum": "111111"}
{"SSN":"000000","twitterHandle":"foo"}

如何使用 Kafka Streams 实现这一目标? 我可以从输入主题创建一个 KStream 并将其转换为 KTable 以获取输出主题吗?

更新: 该主题包含来自多个不同用户的事件。用户标识符(SSN、twitterHandle)不固定。用户可能还有其他 id

【问题讨论】:

  • 目前还不清楚您希望在单个输入主题中聚合多少事件。输入主题是否可以包含两个以上的用户事件?
  • 是的,该主题包含针对不同用户的多个事件,我需要根据复合键(标识符 SSN、Twitter 句柄等)将同一用户的事件聚合在一起。一个事件可能只包含一个标识符,例如该用户的 SSN
  • 如何找出哪个 twitterHandle 需要映射到哪个 SSN?
  • @nikitap 最初是没有办法的。但是当事件 3 出现时,它会链接这些键,就像一个复合键。因此,对于事件 1 和 2,我想假设他们是不同的用户,只有在事件 3 出现后,我才会将它们聚合为单个用户

标签: apache-kafka apache-kafka-streams


【解决方案1】:

如果你一味的想要移除消息 1 & 2 并保留消息 3,你可以使用消费者拦截器。

拦截器会盲目地解析 json 消息,检查消息是否同时存在键(并且不为空),然后成功发送消息,否则不发送。在这种情况下,您不需要 kstream 应用程序。消费消息时只需要使用一个拦截器类。

但是,如果您只想拼接 1 和 2 而它们之间没有任何公共密钥,我认为这是不可能的,因为我们不知道哪个 SSN 需要与哪个 twitter 句柄合并。

如果我能以其他方式提供帮助,请告诉我。

【讨论】:

  • 我需要将 1 和 2 发送到下游主题(可能是 KTable),我需要通过 Kafka Connect 输出到 Elasticsearch。然后某个时候,事件 3 到达。我想和 3 一起加入 1&2。下游现在只有 1 个事件。也许我需要以某种方式将 1 和 2 石化。
猜你喜欢
  • 2013-03-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-09-15
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多