【问题标题】:Incorrect 'key' value in Map transformMap 变换中的“键”值不正确
【发布时间】:2020-08-15 04:18:34
【问题描述】:

apache-beam==2.23.0Python 3.8.5DirectRunner

在我的 Map 转换中,我试图为每个元组元素提取 Key 值(在上游 GroupByKey 转换之后)。但输出始终是字符串 'KeyParam' 而不是实际键值

这是最少的代码:

管道代码

p| beam.Create([("2","elem2.1"),("1","elem1.1"),("1","elem1.2")]) \
|"group" >>beam.GroupByKey() \
| "log_PCollection_AfterGrouped" >> beam.Map(myRawProcessor.myReader) \

地图转换代码

class myRawProcessor():
    @classmethod
    def myReader(self,e,
              timestamp=beam.DoFn.TimestampParam,
              window=beam.DoFn.WindowParam,
              watermark=beam.DoFn.WatermarkEstimatorParam,
              key=beam.DoFn.KeyParam,
               *args, **kwargs):
        print("=== === ===")
        print(e)
        print(key)
        return e

输出

> === === === 
> ('2', ['elem1.1']) 
> KeyParam -----> EXPECTED :: '2'
> === === === 
> ('1', ['elem1.2', 'elem1.3']) 
> KeyParam ----> EXPECTED :: '1'

【问题讨论】:

    标签: python apache-beam


    【解决方案1】:

    这是一个错误,请参阅BEAM-10780。同时,避免在这种情况下使用DoFn.KeyParam

    【讨论】:

    • 非常感谢!奇怪的是,如果我在函数参数中删除 watermark=beam.DoFn.WatermarkEstimatorParam,它工作正常......
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-11-14
    相关资源
    最近更新 更多