【发布时间】: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<Long, WindData>和KV<Long, TempeData>,其中的键是时间戳在时间戳属于的两种情况下的 bin (例如:2200 属于密钥 2000,因此您必须向下舍入到千位)。创建组后,您可以选择最小值或最大值或任何您需要的传感器值。希望这会有所帮助:)
标签: google-cloud-platform google-cloud-dataflow apache-beam google-cloud-pubsub