【问题标题】:Apache Beam - How to associate transformed record with original?Apache Beam - 如何将转换后的记录与原始记录相关联?
【发布时间】:2022-06-15 04:41:03
【问题描述】:

我正在使用 Python SDK 创建一个 Apache Beam 管道,以从 PubSub 读取数据并写入 BigQuery。我正在尝试保留来自 PubSub 的原始消息,以便如果有任何错误,我可以写出要修复的原始记录,然后重新处理。我完成这项工作的最简单方法是使用包含原始消息和工作消息的元组:

(initial_message, working_message)

然后,当我进行 Map 转换时,我转换工作消息并返回元组,保持原始消息不变:

pipeline = (pipeline
            | 'Read PubSub' >> ReadFromPubSub(subscription=subscription)
            | 'Decode' >> Map(lambda record: (record, record.decode('utf-8')))
            | 'Example Transform' >> Map(lambda record: (record[0], some_transformation(record[1])))
)

在写入 BigQuery 之前,这似乎很有效:

write_results = (
    pipeline
    | 'Extract working message' >> Map(lambda record: record[1])
    | 'Write to BigQuery' >> WriteToBigQuery(table=table,
                                             project=project,
                                             schema=schema,
                                             create_disposition=create_disposition,
                                             write_disposition=write_disposition,
                                             insert_retry_strategy=insert_retry_strategy
)

write_results['FailedRows'] | 'Handle write failures' >> ?

然后如何将失败的行与原始消息相关联?

【问题讨论】:

    标签: python google-bigquery streaming apache-beam google-cloud-pubsub


    【解决方案1】:

    感谢您的提问。有没有办法为每条记录分配标识符?例如,您可以为每条记录添加一个唯一标识符,如下所示:

    def assign_id(input_record):
      return RecordWithId(id=uuid.uuid4(),  # Generate a random unique ID for it
                          record=input_record)
    
    def append_id(record_with_id):
      record_with_id.record['_beam_id'] = record_with_id.id
    
    data = (pipeline
                | 'Read PubSub' >> ReadFromPubSub(subscription=subscription)
                | 'AssignId' >> Map(lambda record: assign_id(record))
                | 'Decode' >> Map(lambda record: (record, RecordWithId(record.id, record.record.decode('utf-8'))))
                | 'Append Id to Row' >> Map(lambda pair: (pair[0], append_id(pair[1]))
                | 'Example Transform' >> Map(lambda record: (record[0], some_transformation(record[1])))
    )
    
    write_results = .... # Write to BQ
    
    # And finally, you would do:
    
    kv_failures = write_results['FailedRows'] | KeyBy(lambda row: row['_beam_id'))
    kv_original = data | KeyBy(lambda row_w_id: row_w_id.id)
    
    joined_data = (kv_failures, kv_original) | CoGroupByKey()
    

    这有意义吗?然后您可以处理joined_data

    【讨论】:

      猜你喜欢
      • 2017-09-23
      • 1970-01-01
      • 2020-01-29
      • 1970-01-01
      • 1970-01-01
      • 2015-04-25
      • 1970-01-01
      • 1970-01-01
      • 2017-02-21
      相关资源
      最近更新 更多