【问题标题】:Flink: Append an event to the end of finite DataStreamFlink:将事件附加到有限数据流的末尾
【发布时间】:2019-07-08 07:13:25
【问题描述】:

假设有一个带有事件的有限 DataStream(例如来自数据库源)

  • a1, a2, ..., an

如何在此流中再追加一个事件b 以获取

  • a1, a2, ..., an, b

(即在所有原始事件之后输出添加的事件,保留原始顺序)?

我知道所有有限流在所有事件之后都会发出MAX_WATERMARK。那么,有没有办法“捕捉”这个水印并在它之后输出附加事件?

(不幸的是,.union()ing 原始 DataStream 与另一个 DataStream 组成的单个事件(时间戳设置为 Long.MaxValue)然后使用 this answer 对联合流进行排序不起作用。)

【问题讨论】:

  • 你提前知道计数吗?还有,如果是有限集,为什么不能用DataSet API代替DataStream呢?

标签: apache-flink flink-streaming


【解决方案1】:

也许我遗漏了一些东西,但似乎您可以简单地拥有一个 ProcessFunction 并为遥远的将来某个地方设置一个事件时间计时器,以便它仅在 MAX_WATERMARK 到达时触发。然后在 onTimer 方法中,如果 currentWatermark 为 MAX_WATERMARK,则发出该特殊事件。

【讨论】:

  • 谢谢!它对我有用。我不知道只有当水印到达时 onTimer 才会触发。
【解决方案2】:

另一种方法可能是将原始数据源“包装”在另一个数据源中,当委托对象的run() 方法返回时,它会发出最终元素。当然,您需要小心调用所有委托方法。

【讨论】:

    猜你喜欢
    • 2014-09-09
    • 2015-04-09
    • 1970-01-01
    • 2017-03-26
    • 2017-06-15
    • 1970-01-01
    • 1970-01-01
    • 2012-06-28
    • 2015-01-25
    相关资源
    最近更新 更多