【问题标题】:Merge two PCollection (Apache beam)合并两个 PCollection(Apache 梁)
【发布时间】:2020-01-10 08:31:02
【问题描述】:

我在云存储中有两个文件。包含 Avro 格式的 File1,其中包含来自温度传感器的数据。

time_stamp     |  Temperature
1000           |  T1
2000           |  T2
3000           |  T3
4000           |  T3
5000           |  T4
6000           |  T5

包含来自风传感器的数据的 Avro 格式的 File2。

time_stamp     |  wind_speed
500            |  w1
1200           |  w2
1500           |  w3
2200           |  w4
2500           |  w5
3000           |  w6

我想像下面这样组合输出

time_stamp |Temperature|wind_speed
1000       |T1         |w1 (last earliest reading from wind sensor at 500)
2000       |T2         |w3 (last earliest reading from wind sensor at 1500)
3000       |T3         |w6 (wind sensor reading at 3000)
4000       |T3         |w6 (last earliest reading from wind sensor at 3000)
5000       |T4         |w6 (last earliest reading from wind sensor at 3000)
6000       |T5         |w6(last earliest reading from wind sensor at 3000)

我正在寻找 apache 梁中的解决方案来合并上述文件。现在它正在从文件中读取,但将来它可能会通过 pubsub 来。我想找出组合两个 PCollection 并创建另一个 PCollection tempDataWithWindSpeed 的自定义方式。

     PCollection<Temperature> tempData = p.apply(AvroIO
         .read(AvroAutoGenClass.class)
         .from("gs://my_bucket/path/to/temp-sensor-data.avro")

     PCollection<WindSpeed> windData = p.apply(AvroIO
         .read(AvroAutoGenClass.class)
         .from("gs://my_bucket/path/to/wind-sensor-data.avro")

     PCollection<WindSpeed> tempDataWithWindSpeed = ?

【问题讨论】:

  • 有几种解决方案。你能添加更多细节吗?例如,温度时间戳是否与示例中显示的一样有规律?是流处理还是批处理?在管道中合并后,您是否执行了许多额外的转换?什么样的变换?
  • 这是一个很好的例子,如何加入他们:beam.apache.org/documentation/pipelines/design-your-pipeline/…
  • @guillaumeblaquiere 转换如何影响解决方案。现在是批处理。
  • @jszule 我看到了那个例子,使用的连接键是用户名。我没有直接加入密钥,我需要一些自定义解决方案才能加入。
  • 你仍然可以加入源,只需要创建一个 KV 值,如 KV&lt;Long, WindData&gt;KV&lt;Long, TempeData&gt;,其中的键是时间戳在时间戳属于的两种情况下的 bin (例如:2200 属于密钥 2000,因此您必须向下舍入到千位)。创建组后,您可以选择最小值或最大值或任何您需要的传感器值。希望这会有所帮助:)

标签: google-cloud-platform google-cloud-dataflow apache-beam google-cloud-pubsub


【解决方案1】:

@jszule 的评论通常是 Dataflow/Beam 的一个很好的答案:最受支持的连接是当两个 PCollection 有一个公共键时。对于大多数数据,Beam 可以找出一个模式,您可以使用CoGroup.join 转换。您必须做出的设计决策是如何选择键,例如向下舍入到最接近的 1000。

您的用例有一个复杂之处:您需要在时间序列中为没有数据的键结转值。解决方案是使用状态和计时器来生成“缺失”值。您仍然需要仔细选择键,因为状态和计时器是每个键和窗口的。状态和计时器也可以在批处理模式下工作,因此这是一个批处理/流式统一解决方案。

您可能想阅读有关该主题的this blog post by Reza Rokni and myself,或this talk by Reza at the Beam Summit Berlin 2019

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-03-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-07
    • 1970-01-01
    • 2013-08-13
    相关资源
    最近更新 更多