【问题标题】:Flink Tumbling Window labellingFlink 翻滚窗口标注
【发布时间】:2019-03-07 15:25:00
【问题描述】:

我有一个 flink 应用程序接收以下格式的数据流的场景:

{ "event_id": "c1s2s34", "event_create_timestamp": "2019-03-07 11:11:23", "amount": "104.67" }

我正在使用以下翻转窗口来查找过去 60 秒内输入流的总和、计数和平均量。

keyValue.timeWindow(Time.seconds(60))

但是如何标记聚合结果,以便我可以说 16:20 到 16:21 之间的输出数据流聚合结果是 sum x、count y 和 average z。

任何帮助都会被占用。

【问题讨论】:

  • 您希望如何使用结果——您是要打印它们,还是将它们写入文件,或者将它们发送到 Kafka,...?
  • 嗨,David,我想将结果发送到 Kinesis Firehose。

标签: apache-flink data-stream


【解决方案1】:

如果您查看 Flink 培训站点中的窗口示例 -- https://training.ververica.com/exercises/hourlyTips.html -- 您将看到如何使用 ProcessWindowFunction 从包含时间信息等的窗口创建输出事件的示例。想法是 ProcessWindowFunction 上的 process() 方法传递一个 Context ,该 Context 又包含 Window 对象,您可以从中确定窗口的开始和结束时间,例如,context.window().getEnd()

然后您可以安排您的 ProcessWindowFunction 返回元组或 POJO,其中包含您希望包含在报告中的所有信息。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-08-03
    • 2017-05-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多