【发布时间】: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
我想要的是提供一些额外的参数arg1 和kwarg1,它们是内部其他转换所需的。但我不知道这是否是一种有效的方式,也不知道如何在管道中使用它。
谁能给我指个方向?
【问题讨论】:
标签: python apache-beam