You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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+的特性,简洁易维护,适配这类结构化数据场景:

  1. 定义对应的数据类:
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
  1. 在流水线中添加转换步骤:
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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 13:22:07