如何在Apache Beam流水线中解析DB读取的PCollection并传递至后续步骤
问题
我用Apache Beam写了一个小型流水线,通过beam-postgres输入连接器从数据库表生成PCollection。目前读取的结果是字典格式(示例输出见下文),想把它解析成对象,传递给流水线后续步骤并访问其属性,该怎么实现?
代码示例
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from psycopg.rows import dict_row from beam_postgres.io import ReadAllFromPostgres def __trigger_bill_fetch_job(self): print("Triggering bill-fetch job") pipeline = beam.Pipeline() read_from_db = ReadAllFromPostgres( "host={host} dbname={dbName} user={user} password={password}", "SELECT * FROM comparison_bill_data_requests WHERE status='PENDING' AND bill_event_received=true and bill_detail_event_received=true", dict_row, ) result = pipeline | "ReadPendingRecordsFromDB" >> read_from_db | "Print result" >> beam.Map(print) pipeline.run().wait_until_finish() print("read_from_db done", read_from_db)
输出示例
{'id': '1', 'bill_id': 'bill-1', 'account_id': None, 'bill_event_received': True, 'bill_detail_event_received': True, 'status': 'PENDING', 'commodity_type': None, 'bill_start_date': datetime.datetime(2023, 12, 6, 11, 52, 28, 78945), 'bill_end_date': datetime.datetime(2023, 12, 6, 11, 52, 28, 78945), 'tenant_id': 'tenant-1', 'created_at': datetime.datetime(2023, 12, 6, 11, 52, 28, 78945), 'updated_at': datetime.datetime(2023, 12, 6, 11, 53, 21, 300224)}
解决方案
可以通过定义数据类、普通类或命名元组将字典转换为对象,再通过Beam的Map操作完成转换,后续步骤就能直接访问对象属性。
方法1:使用dataclass(推荐)
dataclass是Python 3.7+的特性,简洁易维护,适配这类结构化数据场景:
- 定义对应的数据类:
from dataclasses import dataclass import datetime @dataclass class BillDataRequest: id: str bill_id: str account_id: str | None bill_event_received: bool bill_detail_event_received: bool status: str commodity_type: str | None bill_start_date: datetime.datetime bill_end_date: datetime.datetime tenant_id: str created_at: datetime.datetime updated_at: datetime.datetime
- 在流水线中添加转换步骤:
def dict_to_bill_request(record: dict) -> BillDataRequest: return BillDataRequest(**record) # 修改流水线逻辑 result = ( pipeline | "ReadPendingRecordsFromDB" >> read_from_db | "ConvertToObject" >> beam.Map(dict_to_bill_request) | "PrintObjectProperties" >> beam.Map(lambda obj: print(f"ID: {obj.id}, 账单ID: {obj.bill_id}")) )
方法2:使用普通类
如果需要自定义初始化逻辑,可使用普通类:
import datetime class BillDataRequest: def __init__(self, id, bill_id, account_id, bill_event_received, bill_detail_event_received, status, commodity_type, bill_start_date, bill_end_date, tenant_id, created_at, updated_at): self.id = id self.bill_id = bill_id self.account_id = account_id self.bill_event_received = bill_event_received self.bill_detail_event_received = bill_detail_event_received self.status = status self.commodity_type = commodity_type self.bill_start_date = bill_start_date self.bill_end_date = bill_end_date self.tenant_id = tenant_id self.created_at = created_at self.updated_at = updated_at # 转换函数 def dict_to_bill_request(record: dict) -> BillDataRequest: return BillDataRequest(**record) # 流水线使用方式与方法1一致
方法3:使用namedtuple
如果仅需不可变的简单对象,可使用namedtuple:
from collections import namedtuple import datetime BillDataRequest = namedtuple('BillDataRequest', [ 'id', 'bill_id', 'account_id', 'bill_event_received', 'bill_detail_event_received', 'status', 'commodity_type', 'bill_start_date', 'bill_end_date', 'tenant_id', 'created_at', 'updated_at' ]) # 转换函数直接解包字典 def dict_to_bill_request(record: dict) -> BillDataRequest: return BillDataRequest(**record)
注意事项
- 确保字典键与类/namedtuple的字段名完全匹配,否则会触发参数错误;字段名不匹配时,需在转换函数中手动映射。
- 处理
None值时,确保类字段类型支持(比如用str | None标注可选类型)。 - 若数据库返回字段类型与类定义不一致,需在转换函数中做类型转换(示例中日期已为
datetime类型,无需额外处理)。
内容的提问来源于stack exchange,提问作者Sumit Desai
相关产品推荐
相关产品推荐

