【问题标题】:How to extract contents from PCollection in Cloud Dataflow?如何从 Cloud Dataflow 中的 PCollection 中提取内容?
【发布时间】:2015-01-18 23:15:50
【问题描述】:

只是想知道如何从 PCollection 中提取内容? 假设我已经应用了 Count.Globally,所以在生成的 PCollection 中只有一个数字,但是如何将其提取为 Long 值?

谢谢。

【问题讨论】:

    标签: google-cloud-dataflow


    【解决方案1】:

    这取决于您希望如何使用该值。

    如果您想在管道完成后读取该值,您可以使用一种写入转换(例如AvroIO.Write)将其写入某个输出,然后您可以从管道完成后执行的任何代码中读取。

    如果您想在管道的后续部分中使用该值,则可以应用 View 转换来生成 PCollectionView,然后您可以将其作为侧面输入传递给其他转换。

    考虑一个简单的例子,目标是打印出计数。直到管道运行后,计数才可用。所以在这种情况下,我们可以执行以下操作

    • 定义一个 DoFn 我们将其应用于计数,以便将 Long 转换为我们要打印的消息。
    • 应用 TextIO.Write 转换将消息写入文件。
    • 运行作业并等待它完成。如果我们想使用 Dataflow Service 执行,我们可以使用 BlockingDataflowRunner 等待作业完成。
    • 作业完成后,读取为获取消息而创建的文本文件并将其打印出来。

    【讨论】:

    • 我想我想得到一个 Long 值并在程序中使用。说在 if 语句中使用。
    • 我添加了一个示例,这是否回答了您的问题?重要的是 PCollection 在管道运行之前不会实现。因此,您需要运行管道才能访问该值。管道也可以(取决于运行程序)异步运行,独立于您的主程序。
    • Jeremy,真的没有直接的方法可以访问计数,即使我们已经运行了管道?将其写入磁盘然后再读回似乎真的很间接。
    • 您可以在管道中实现它,但不能在主程序中实现。您的主程序正在构建转换图。如果您在 Dataflow 服务上执行该图,则计算实际上是在与您的主程序运行的机器不同的机器上执行的。因此,为了让您的主程序读取数据,您需要使用 IO 转换将其写入您的主程序可访问的某个位置(例如 File/BigQuery/Datastore/Pubsub)。但是,如果您想在以后的转换中读取它,Dataflow 将为您传递该值。
    【解决方案2】:

    您必须始终将PCollection 视为。您应用了每个窗口创建单个值的转换这一事实并不能保证实际上只有单个值。这取决于窗口策略 - 因此在您使用 GlobalWindow 的情况下可能会有单个值,但对于其他类型的窗口函数(例如滑动窗口)会有很多值。

    因此,无法直接提取此单个值(例如 PCollection.get() 之类的东西) - 返回值必须是流。如果要从 PCollection 检索结果,则必须对其应用转换,将其存储在某处。有一组丰富的内置 IO 模块(参见here)。如果您想检索结果值,并稍后在程序中使用它,最好的选择是将其存储在您选择的某个共享数据库中,并在管道完成后检索此值。请注意,这意味着您的管道是有界(例如批处理,而不是流式传输),否则它将永远不会完成。但是您的问题表明您想到的是有界管道。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-02-09
      • 1970-01-01
      • 2018-06-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多