基于PubSub消息读取BigQuery表时部署模板遇TypeError问题求助
问题根源
你这段代码里的tablename是个PCollection(Beam里的数据集对象),但beam.io.ReadFromBigQuery的table参数只接受字符串格式的表名或者BigQueryTableReference对象,根本不能直接传PCollection,这就是为啥会报TypeError: expected string or buffer错误。
解决方法
因为PubSub消息是运行时才会产生的,没法在Pipeline初始化阶段提前确定要读的BigQuery表名,所以不能用ReadFromBigQuery作为起始读取步骤。正确的思路是先接收PubSub的消息,然后在ParDo里针对每个表名,用BigQuery客户端动态读取对应表的数据:
import json import apache_beam as beam from google.cloud import bigquery class ReadBigQueryTable(beam.DoFn): def process(self, tablename): # 初始化BigQuery客户端 client = bigquery.Client() # 构造查询语句读取整张表 query = f"SELECT * FROM `{tablename}`" query_job = client.query(query) # 将查询结果转成字典格式输出 for row in query_job.result(): yield dict(row) subscription = 'my_subscription' with beam.Pipeline() as p: table_names = (p | '接收PubSub消息' >> beam.io.ReadFromPubSub(subscription=subscription) | '转成字典' >> beam.Map(lambda x: json.loads(x)) | '提取表名' >> beam.Map(lambda x: x['tablename']) ) bq_data = (table_names | '读取对应BigQuery表' >> beam.ParDo(ReadBigQueryTable()) | # 这里加后续的转换操作... )
额外注意
- 要保证运行这个Pipeline的服务账号,同时拥有PubSub订阅权限和BigQuery数据读取权限
- 如果目标表的数据量很大,建议在ParDo里做分页处理,避免内存占用过高导致异常
内容的提问来源于stack exchange,提问作者FrancoisDuCoq
相关产品推荐
相关产品推荐

