【发布时间】:2021-09-06 12:47:15
【问题描述】:
我已经阅读了很多教程,他们已经解释过 transform 的输出是 apache Beam 中的 Pcollection。 谁能解释一下,Pcollection 是如何存储的,如果我们应用任何转换,它会返回什么数据类型。 是python字典、元组、列表吗?
【问题讨论】:
标签: apache-spark pyspark apache-beam dataflow
我已经阅读了很多教程,他们已经解释过 transform 的输出是 apache Beam 中的 Pcollection。 谁能解释一下,Pcollection 是如何存储的,如果我们应用任何转换,它会返回什么数据类型。 是python字典、元组、列表吗?
【问题讨论】:
标签: apache-spark pyspark apache-beam dataflow
Apache Beam 以 延迟 方式执行管道。这意味着您在构建管道时无法访问 PCollection 的元素 - 因为尚未计算 PCollection:
p = beam.Pipeline(runner='....')
input_pcollection = p | ReadFromText(...)
result_pcollection = input_pcollection | beam.Filter(...) | beam.Combine(...)
此时,您只告诉 Beam 应该执行哪些操作 未来 - 但这些操作尚未执行。
要实际计算 PCollection,您必须:
p.run().wait_until_finish()
要实际检查 PCollection,请尝试在 Jupyter 或 Collab 笔记本中以交互方式运行:https://colab.sandbox.google.com/github/apache/beam/blob/master/examples/notebooks/get-started/try-apache-beam-py.ipynb
【讨论】: