无法将GCS存储桶中的CSV加载到BigQuery表 仅前两列有值其余为NULL
问题排查与解决方案
错误根因
你遇到的仅前两列有数据、剩余列全为NULL的问题,主要由以下几个代码问题导致:
- CSV解析参数缺失:airbnb数据集的
name、host_name等文本字段常包含逗号,且这类字段默认用双引号包裹,你调用csv.reader时未指定quotechar='"',导致带逗号的文本字段被错误拆分为多个字段,整行字段顺序错位,最终只有前两个拆分结果能匹配字典的前两个key,后续字段全部匹配失败。 - 代码缩进错误:
parse_file函数下的for循环没有正确缩进,函数无法返回完整的字段列表,部分场景下只会返回前两个拆分片段。- 流水线的转换、写入步骤缩进混乱,部分步骤没有被包含在
with beam.Pipeline()的上下文块中,导致流水线执行逻辑不完整。
- 缺少字段校验逻辑:没有校验解析后得到的字段数量是否和预期的16个字段匹配,字段错位后无法及时发现拦截。
正确实现方案
以下是可直接运行的DataFlow流水线实现,包含完整的字段解析、类型转换、异常过滤逻辑:
import csv import apache_beam as beam from apache_beam.io import ReadFromText from apache_beam.options.pipeline_options import PipelineOptions # 数据集对应字段列表 CSV_COLUMNS = ('id', 'name', 'host_id', 'host_name', 'neighbourhood_group', 'neighbourhood', 'latitude', 'longitude', 'room_type', 'price', 'minimum_nights', 'number_of_reviews', 'last_review', 'reviews_per_month', 'calculated_host_listings_count', 'availability_365') def parse_file(element): # 配置quotechar处理带引号、逗号的文本字段 for line in csv.reader([element], delimiter=',', quotechar='"', skipinitialspace=True): # 校验字段数量,过滤错位行 if len(line) == len(CSV_COLUMNS): return line return None def format_to_bq_row(values): if not values: return None row = dict(zip(CSV_COLUMNS, values)) # 字段类型转换,避免类型不匹配写入失败 try: row['id'] = int(row['id']) row['host_id'] = int(row['host_id']) row['latitude'] = float(row['latitude']) row['longitude'] = float(row['longitude']) row['price'] = int(row['price']) row['minimum_nights'] = int(row['minimum_nights']) row['number_of_reviews'] = int(row['number_of_reviews']) row['reviews_per_month'] = float(row['reviews_per_month']) if row['reviews_per_month'] else None row['calculated_host_listings_count'] = int(row['calculated_host_listings_count']) row['availability_365'] = int(row['availability_365']) except ValueError: # 过滤类型转换失败的异常行 return None return row def run(input_path, table_spec, pipeline_args=None): pipeline_options = PipelineOptions(pipeline_args, save_main_session=True) # 所有流水线步骤都放在with上下文块内 with beam.Pipeline(options=pipeline_options) as p: ( p | '读取GCS CSV文件' >> ReadFromText(input_path, skip_header_lines=1) | '解析CSV行' >> beam.Map(parse_file) | '过滤无效行' >> beam.Filter(lambda x: x is not None) | '转换为BQ表结构' >> beam.Map(format_to_bq_row) | '过滤格式异常行' >> beam.Filter(lambda x: x is not None) | '写入BigQuery' >> beam.io.WriteToBigQuery( table_spec, schema=""" id:INTEGER, name:STRING, host_id:INTEGER, host_name:STRING, neighbourhood_group:STRING, neighbourhood:STRING, latitude:FLOAT, longitude:FLOAT, room_type:STRING, price:INTEGER, minimum_nights:INTEGER, number_of_reviews:INTEGER, last_review:DATE, reviews_per_month:FLOAT, calculated_host_listings_count:INTEGER, availability_365:INTEGER """, write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) ) if __name__ == '__main__': # 替换为你的实际路径 input_path = 'gs://你的存储桶名/AB_NYC_2019.csv' table_spec = '你的GCP项目ID:你的BigQuery数据集名.airbnb_nyc' # 运行时传入DataFlow参数即可切换到DataFlow运行 run(input_path, table_spec)
补充说明
- 如果不需要自定义字段处理逻辑,全量同步CSV到BigQuery也可以直接使用BigQuery原生LOAD JOB,无需运行DataFlow流水线,成本更低效率更高。
- 运行DataFlow作业时,只需在调用脚本时传入
--runner DataflowRunner --project 你的项目ID --region 运行区域 --temp_location gs://你的临时存储桶/tmp参数即可。
内容的提问来源于stack exchange,提问作者Akhil Kv
相关产品推荐
相关产品推荐

