【问题标题】:dataflow streaming job early results数据流流作业早期结果
【发布时间】:2020-09-01 00:50:32
【问题描述】:

相关Early results from GroupByKey transform

  1. 以流模式从 GCS 读取 avro 源文件
  2. 过滤实验事件和输出键值对。 键->“experimentId”:“aa”,“experimentVariant”:2,“uuid”:abbcd 值-> 事件日期
  3. 使用丢弃窗格修复了缓冲元素的窗口(60-120 秒)
  4. 组合每个键以收集不同的日期。 输出示例 -> 一窗结果 键->“experimentId”:“aa”,“experimentVariant”:2,“uuid”:abbcd 值 -> 设置(“2020-06-01”,“2020-06-02”) 下一个窗口结果 键->“experimentId”:“aa”,“experimentVariant”:2,“uuid”:abbcd 值 -> 设置(“2020-06-03”)
  5. 写入 gcs

问题是即使窗口只有 60 秒,组合步骤也不会长时间提供输出。

【问题讨论】:

    标签: google-cloud-dataflow apache-beam


    【解决方案1】:

    watermark 到达窗口末尾之前,诸如组合之类的聚合步骤不会给出输出(除非设置了非默认触发)。对于基于文件的源,可能会出现记录未按时间戳排序的情况,因此必须先读取整个文件,然后才能安全地推进水印。在 Dataflow 上,您可以see this in the UI

    【讨论】:

    • 我正在使用下面的触发器,而且我有数千个文件,它一直在等待获取所有数据。它应该在读取一些文件后触发输出 atleast.p.apply("f", Window .>>into(FixedWindows.of(Duration.standardSeconds(120))).trigger(AfterEach.inOrder(rtdly.forever(AfterFirst.of(AfterPane.elementCountAtLeast(10000),APT. pastFirstElementInPane().plusDelayOf(Duration.standardSeconds(60)))) .orFinally(AfterWatermark.pastEndOfWindow()),rtdly.forever(APT.pastFirstElementInPane() .plusDelayOf(Duration.standardSeconds(240))))).wal( DR.standardMinutes(10)).dfp())
    • 水印是一个承诺,它永远不会输出时间戳T的元素。直到读取最后一个文件的最后一条记录,水印才能安全前进。 (如果您的所有输入都是文件,听起来您应该将其作为批处理管道运行。)
    • 是否可以在最后一个文件的最后一条记录之前将水印提前?它将保留内存中的所有记录直到那时。实际上用例是流式传输,源文件每 2 分钟添加一次。
    • 如果你使用像 TextIO.ReadAll() 这样的东西,每个文件名都有一个时间戳(文件中的每个元素都会在这个时间戳上发出),并且会随着文件的读取而前进。如果您需要比这更精细的东西,请考虑使用 SDF。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-12-02
    • 1970-01-01
    • 2020-08-05
    • 2021-01-22
    • 2021-11-12
    相关资源
    最近更新 更多