【问题标题】:Aggregating messages from Kafka in a storm bolt based on a key基于密钥在风暴螺栓中聚合来自 Kafka 的消息
【发布时间】:2018-08-22 08:16:27
【问题描述】:

背景

后端日志处理系统已经在 Kafka 和 Storm 集群中到位。

用例

在后端生成并记录多个特定类型的事件X。每个都包含一个 ID,比如 userid。现在这些事件被一个风暴螺栓消耗并提取 useid 和其他一些字段说 userdata 并写入 kafka 中的另一个主题,比如 data 主题。

现在一些其他拓扑从这个data 主题消费。它使用单个userid 和不同的userdata 查找多个此类事件。如果有n 这样的记录向他们展示,则需要采取一些措施。

问题

如何使用来自 kafka 的一些关键数据在风暴螺栓中聚合? 有些用户可能会在 20 分钟内达到N 记录计数,有些可能需要几个小时,具体取决于用户交互,因此事件记录在后端。目标是当此类记录的计数达到某些 N 时,获取所有用户 ID 和相应的 usedata

【问题讨论】:

    标签: apache-kafka apache-storm aggregation


    【解决方案1】:

    这不是 Storm 特有的问题,而是与用户会话管理有关。如果您希望您的系统面临大量需要很长时间才能获得特定状态的会话(在您的情况下达到 n 事件)并最终在此期间建立大量数据,那么您需要考虑到这一点在您的设计中,意味着明智地选择 n 并围绕它构建大量集成测试,以检查您的系统在负载下保持响应。

    你可以

    • 考虑根据负载和统计信息使 n 动态化(我猜这就是 DevOps 的全部意义所在)
    • 只为userid 存储n 并将数据保存在数据库或文件系统中,如果n 达到临界值,则将其取回。这在某种程度上与流式拓扑的想法相矛盾,但它开始变得有意义,特别是如果您需要在处理数据后保留(部分)数据。

    并不是说如果您想部署软件更新,除了n 之外,您还需要考虑时间值t,因为会话越短,持续部署就越容易。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-09-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多