【问题标题】:How to join a realtime stream and a delayed stream by event time如何按事件时间加入实时流和延迟流
【发布时间】:2020-06-09 18:24:37
【问题描述】:

我正在尝试使用 CoProcessFunction 加入 2 个流。输入流A之一是实时生成的。然而,另一个输入流 B 由一个延迟 1 天的每日计划作业加载,这意味着今天放入流中的事件始终具有昨天的事件时间。

话虽如此,流 B 的水印总是比 A 的水印晚约 1 天,所以我想来自 A 的很多事件将被缓冲在内存中。我想知道是否有办法解决问题。一些额外的背景,流 A 和 B 都是 kinesis 流(我使用的是FlinkKinesisConsumer),保留期 = 7 天。

提前致谢!

【问题讨论】:

  • 你能告诉更多/也许提供一个连接看起来如何的例子吗?你想加入来自 B 的时间戳相等/更低的时间戳吗?
  • 当然,举个简单的例子:两个流中的事件都有一个字段dataId,所以对于流B(延迟的那个)中的给定事件X,我想找到0个或多个流 A 中具有相同 dataId 且事件时间不晚于事件 X 的事件时间后 5 分钟的事件。
  • 对机制的期望是什么?你不想要任何缓冲?
  • 我知道未处理的事件将被缓冲在内存中,在我的情况下,来自当天生成的流 A 的事件将被缓冲。我只是担心应用程序可能会耗尽内存并崩溃。

标签: apache-flink flink-streaming


【解决方案1】:

如果您在缓冲多个元素时担心内存,那么您应该看看各种不同的state backends,尤其是 RocksDb。这样,状态将保存在磁盘上而不是内存中。

这应该可以轻松解决您在一天内遇到缓冲元素的问题,因为限制状态大小的唯一因素是可用磁盘空间,这通常很便宜,在大多数情况下不应该成为问题。

【讨论】:

    【解决方案2】:

    从 Flink 1.8.1 开始,Flink Kinesis 消费者支持事件时间对齐(选择性地从拆分中读取,以确保各个消费者在事件时间中均匀推进)。详情请见Event Time Alignment for Shard Consumers

    Flink 社区正在努力为跨源的事件时间同步提供更通用的支持,以便在源之间存在显着事件时间错位的情况下有效地实现事件时间连接。目前唯一的解决方案(除非您使用 Kinesis)是使用 Flink 状态来缓冲前面的流,这可能会导致非常大的检查点和显着的背压。

    通用事件时间对齐的基础正在作为FLIP-27 / FLINK-10740 的一部分实施,之后必须重新处理源以利用这种新机制。

    【讨论】:

    • 感谢您提供的信息!我阅读了事件时间对齐部分,我的理解是我可以控制如何在同一个运动流中的不同分片之间同步多个水印,但是在同步多个运动流之间的消耗率方面,除了在 Flink 中缓冲之外,我还能如何做到这一点状态?
    【解决方案3】:

    我认为您面临的情况实际上可能不是问题,假设您的“快速”流是通过 Kafka 之类的东西进入的,它可以充当缓冲区并保留消息。 (如果不是这种情况,应该直接在加入前写信给Kafka自己创建这种情况)

    虽然我没有测试它,但我希望在 Flink 内部可用的缓冲区已满之前提取快速流。此时,它会简单地停止或减慢摄取,直到第二个慢流进入以“清理”所有等待加入的消息,之后快速流可以再次开始移动。

    请注意,这可能需要两个流中的消息以大致相同的顺序出现。

    【讨论】:

      猜你喜欢
      • 2013-06-12
      • 1970-01-01
      • 2015-09-03
      • 2019-01-04
      • 1970-01-01
      • 1970-01-01
      • 2019-09-29
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多