使用Apache Beam的ReadFromBigQuery方法触发TypeError错误求助
解决用ReadFromBigQuery读取GCS CSV文件触发的TypeError错误
错误原因
beam.io.ReadFromBigQuery是专门读取BigQuery数据表的组件,它的table参数要求传入BigQuery表标识符(格式如项目ID.数据集ID.表ID),而非GCS存储桶的CSV文件路径。你传入GCS路径会导致内部类型校验失败,触发TypeError: isinstance() arg 2 must be a type or tuple of types错误。
正确解决方案
读取GCS上的CSV文件,需使用Apache Beam中处理文本/CSV的专用组件,以下是两种可行实现:
方式1:用ReadFromText读取并解析CSV
import apache_beam as beam import csv from io import StringIO def parse_csv(row): reader = csv.reader(StringIO(row)) return next(reader) def print_row(row): print(row) pipeline = beam.Pipeline() test_pipeline = (pipeline | "读取GCS文件" >> beam.io.ReadFromText('gs://cloud-samples-data/bigquery/sample-transactions/transactions.csv', skip_header_lines=1) | "解析CSV" >> beam.Map(parse_csv) | "打印行数据" >> beam.Map(print_row)) pipeline.run()
skip_header_lines=1用于跳过CSV表头行(根据文件实际情况调整)- 借助
csv模块解析每行内容,避免手动分割字段出现异常
方式2:用Beam DataFrame读取CSV
import apache_beam as beam from apache_beam.dataframe.convert import to_pcollection def print_row(row): print(row) pipeline = beam.Pipeline() # 读取CSV为DataFrame df = beam.dataframe.read_csv('gs://cloud-samples-data/bigquery/sample-transactions/transactions.csv') # 转换为PCollection处理 test_pipeline = (to_pcollection(df) | "打印行数据" >> beam.Map(print_row)) pipeline.run()
- 适合熟悉Pandas操作的场景,代码更简洁
内容的提问来源于stack exchange,提问作者Sahal Nurain
相关产品推荐
相关产品推荐

