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

如何用Apache Beam Python处理Pub/Sub消息并写入BigQuery

Apache Beam Python 实现Pub/Sub消息到BigQuery的转换

核心转换逻辑

将Pub/Sub消息的原生字段映射到BigQuery目标表的指定列:

  • data:把原始二进制消息数据编码为Base64字符串(确保二进制数据可在BigQuery中正常存储)
  • attr:将消息的attributes键值对序列化为JSON字符串
  • key:使用Pub/Sub消息的messageId作为唯一标识
  • publishTime:将消息发布时间转换为ISO格式字符串(适配BigQuery的TIMESTAMP类型)

完整代码示例

import apache_beam as beam
import base64
import json
from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions

def parse_pubsub_message(message):
    # 编码二进制data为Base64字符串
    encoded_data = base64.b64encode(message.data).decode('utf-8')
    
    # 将attributes字典序列化为JSON字符串
    attr_json = json.dumps(message.attributes)
    
    # 提取messageId作为key字段值
    message_key = message.message_id
    
    # 转换发布时间为ISO格式字符串
    publish_time = message.publish_time.isoformat()
    
    # 返回适配BigQuery表结构的字典
    return {
        'data': encoded_data,
        'attr': attr_json,
        'key': message_key,
        'publishTime': publish_time
    }

def run():
    # 配置Pipeline流式处理选项
    pipeline_options = PipelineOptions()
    pipeline_options.view_as(StandardOptions).streaming = True
    
    with beam.Pipeline(options=pipeline_options) as p:
        # 从指定Pub/Sub订阅读取消息
        pubsub_messages = p | 'Read Pub/Sub Messages' >> beam.io.ReadFromPubSub(
            subscription='projects/YOUR_PROJECT/subscriptions/YOUR_SUBSCRIPTION'
        )
        
        # 转换消息格式以匹配BigQuery表结构
        transformed_data = pubsub_messages | 'Transform Message Format' >> beam.Map(parse_pubsub_message)
        
        # 写入BigQuery表
        transformed_data | 'Write to BigQuery' >> beam.io.WriteToBigQuery(
            table='YOUR_PROJECT:YOUR_DATASET.YOUR_TABLE',
            schema='data:STRING, attr:STRING, key:STRING, publishTime:TIMESTAMP',
            write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
            create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED
        )

if __name__ == '__main__':
    run()

关键细节说明

  1. 字段处理:

    • 二进制message.data用Base64编码,避免BigQuery存储二进制数据的兼容性问题
    • message.attributes序列化为JSON字符串,保留完整键值对结构,方便后续解析
    • message.message_id作为key,可用于消息去重或唯一标识
    • message.publish_time转为ISO格式后,BigQuery会自动识别为TIMESTAMP类型
  2. BigQuery表结构:
    若预先创建表,可使用以下DDL:

    CREATE TABLE YOUR_DATASET.YOUR_TABLE (
        data STRING,
        attr STRING,
        key STRING,
        publishTime TIMESTAMP
    )
    
  3. 运行前置要求:

    • 安装依赖:pip install apache-beam[gcp]
    • 配置GCP权限(Pub/Sub读取、BigQuery写入)
    • 替换代码中的项目ID、订阅ID、数据集及表名

内容的提问来源于stack exchange,提问作者Mihir Sharma

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 15:35:36