【问题标题】:How to supply parameters to a composite transform in Apache Beam?如何为 Apache Beam 中的复合变换提供参数?
【发布时间】:2018-12-19 06:36:32
【问题描述】:

我正在使用 Apache Beam 的 Python SDK。

我有几个转换步骤,并希望使它们可重复使用,这让我编写了一个自定义复合转换,如下所示:

class MyCompositeTransform(beam.PTransform):
def expand(self, pcoll, arg1, kwarg1=u'default'):
    result = (pcoll
              | 'Step 1' >> beam.Map(lambda f: SomeFn(f, arg1))
              | 'Last step' >> beam.Map(lambda f: SomeOtherFn(f, kwarg1))
              )
    return result

我想要的是提供一些额外的参数arg1kwarg1,它们是内部其他转换所需的。但我不知道这是否是一种有效的方式,也不知道如何在管道中使用它。

谁能给我指个方向?

【问题讨论】:

    标签: python apache-beam


    【解决方案1】:

    您可以通过PTransform 构造函数提供参数。参数也可以采用侧输入的形式(即从另一个变换输出的数据)。这是一个同时使用“普通”参数和侧输入的示例。

    from typing import Dict, Any, Iterable
    import apache_beam as beam
    
    
    class MyCompositeTransform(beam.PTransform):
    
        def __init__(self, my_arg, my_side_input):
            super().__init__()
            self.my_arg= my_arg
            self.my_side_input= my_side_input
    
        @staticmethod
        def transform(
            element: Dict[str, Any], my_arg: int, my_side_input: Iterable[int]
        ) -> Dict[str, Any]:
            pass
    
        def expand(self, pcoll):
            return pcoll | "MyCompositeTransform" >> beam.Map(
                MyCompositeTransform.transform,
                self.my_arg,
                beam.pvalue.AsIter(self.my_side_input),
            )
    

    使用beam.pvalue 定义侧输入如何传递给转换,例如它是单个值、Iterable 还是具体化为 List

    来自 Beam 的其他示例:(请参阅 PTransformhttps://beam.apache.org/releases/pydoc/2.20.0/_modules/apache_beam/transforms/stats.html

    【讨论】:

    • 很好的答案,应该被接受:它回答了问题,直截了当,它提供了一个清晰的例子,它还提供了关于如何使用光束的有用的侧面输入(双关语)。 pvalue 在 DoFn 中正确使用数据。
    【解决方案2】:

    一般来说,您不能如您所描述的那样在运行时动态地将附加参数传递给转换。当您运行构建管道的控制器程序时,管道的结构被序列化、发送,然后在一群无权访问您的控制器程序的工作人员上并行执行,他们只获得结构和实际您的ParDos 的代码。

    动态参数化执行的一种方法是将额外数据作为额外输入提供,例如创建另一个填充参数值的PCollection,然后将其与主PCollection 连接起来。例如使用side-inputs,或CoGroupByKey

    如果您正在查看 Cloud Dataflow,那么您可能会考虑使用管道 templates with ValueProviders,但不确定它们是否在 pyton 或非 Dataflow 运行器中可用。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-05-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-12-14
      • 1970-01-01
      • 2020-10-17
      相关资源
      最近更新 更多