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

无法将GCS存储桶中的CSV加载到BigQuery表 仅前两列有值其余为NULL

问题排查与解决方案

错误根因

你遇到的仅前两列有数据、剩余列全为NULL的问题,主要由以下几个代码问题导致:

  • CSV解析参数缺失:airbnb数据集的name、host_name等文本字段常包含逗号,且这类字段默认用双引号包裹,你调用csv.reader时未指定quotechar='"',导致带逗号的文本字段被错误拆分为多个字段,整行字段顺序错位,最终只有前两个拆分结果能匹配字典的前两个key,后续字段全部匹配失败。
  • 代码缩进错误:
    1. parse_file函数下的for循环没有正确缩进,函数无法返回完整的字段列表,部分场景下只会返回前两个拆分片段。
    2. 流水线的转换、写入步骤缩进混乱,部分步骤没有被包含在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 10:57:03