【问题标题】:How to merge two kafka topics data when second topic data arrives late当第二个主题数据迟到时如何合并两个kafka主题数据
【发布时间】:2020-08-17 17:16:56
【问题描述】:

有两个来自 mysql 数据库的不同表的 kafka 主题。

表 1 - 交易数据

Table2 - 交易详情数据

现在我需要合并来自这两个 kafka 主题(又名 mysql 表)的数据,并将其作为一个文档推送到 Mongo Db。

虽然我可以使用 kafka 流来做同样的事情,但需要建议如何处理以下情况

案例 1 - Table1 数据到达但 Table2 数据未到达时

情况 2 - Table2 数据到达但 Table1 数据未到达时

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    将您的数据临时存储在窗口键值存储中。

    当stream1的数据到达时:查看stream2的匹配数据是否可用。如果是这样,将数据合并并存储在 MongoDB 中。如果没有,则将流 1 的数据存储在窗口存储中。

    当数据从stream2 到达时:查看stream1 的匹配数据是否可用。如果是这样,组合数据并将其存储在 MongoDB 中。如果没有,则将 stream2 中的数据存储在窗口存储中。

    KafkaStreams 中窗口存储的默认实现是每个分区一个 RocksDB 实例。要完成这项工作,您必须确保两个流具有相同的分区。

    这正是 kafka 流在 KStream.join(Kstream, ...) 后面所做的:

    KStream<String, String> joined = left.join(right,
        (leftValue, rightValue) -> combine(leftValue, rightValue),
        JoinWindows.of(...),
        Joined.with(...)
    );
    

    连接窗口的大小通常是有限的,以避免数据无限长。限制应该是来自不同流的数据的不同到达时间之间的最大差异。

    【讨论】:

      猜你喜欢
      • 2021-08-20
      • 2020-01-07
      • 2021-06-11
      • 1970-01-01
      • 2022-08-11
      • 2017-01-08
      • 1970-01-01
      • 2019-12-05
      • 2019-06-26
      相关资源
      最近更新 更多