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

如何通过Airflow将JSON文件导入BigQuery?报错求助

问题描述

尝试通过Airflow的DAG将GCS存储桶中的JSON文件导入BigQuery,设置每次上传新文件时截断目标表,但执行任务时报错。相关信息如下:

1. JSON文件示例

{"ID":"4238382","Title":"El clon Capítulo 3","Description":"Said y Ali llegan a un acuerdo. Leonidas sale con Yvete y Diogo. Edvaldo no quiere hacerse los exámenes. Jade se rehúsa a usar velo. Lucas se disculpa con Ali. Albieri dice que Ali fue duro con Jade, Ali lo acusa de querer experimentar con humanos.","Program":"El Clon","Season":"1","Episode":"3","Source":"GLOBO TV INTERNACIONAL","Category":"Drama","Syndicator":"CSv2","[CSv2] external_id":"ELCL100002002","[CSv2] pub_win_US_begin":"1661842800","[CSv2] pub_win_US_end":"1754625600","[CSv2] language":"es","[CSv2] title":"El clon Capítulo 3","[CSv2] descriptive_title":"Acuerdo de matrimonio","[CSv2] description":"Said y Ali llegan a un acuerdo. Leonidas sale con Yvete y Diogo. Edvaldo no quiere hacerse los exámenes. Jade se rehúsa a usar velo. Lucas se disculpa con Ali. Albieri dice que Ali fue duro con Jade, Ali lo acusa de querer experimentar con humanos.","[CSv2] supplier":"GLOBO TV INTERNACIONAL","[CSv2] categories":"Drama","[CSv2] rating":"TV-14","[CSv2] subratings":"D","[CSv2] program_type":"NOVELA","[CSv2] entity":"","[CSv2] exception_countries":"US ,	UM ,PR ,	MX ,	AR ,	CL ,	CO ,	PE ,	EC ,	CR ,	SV ,	HN ,	BO ,	PA ,	DO ,	NI ,	PY ,	VE ,	UY ,	GT","[CSv2] episode_type":"","TMS ID":null,"external_id":"ELCL100002002","Content Type":"Entertainment","Release Year":"2001","sports_event_ID":""}

2. 前置准备

已在BigQuery中创建包含所有JSON键作为列的表

3. DAG代码

def source_exists(ti):
    source_folder = list_files_in_bucket(mcp_bucket, mcp_source_folder)
    source_files = []
    for file in source_folder:
        if file.endswith(".json"):
            source_files.append(file)
    if len(source_files) == 0:
        ti.xcom_push(key="json_error_file", value=True)
        raise AirflowException("Source files dont exist")
    ti.xcom_push(key="json_error_file", value=False)

with DAG(
    dag_id=dag_id,
    #schedule_interval=schedule_mcp,
    default_args=default_dag_args,
) as dag:
    source_folder = list_files_in_bucket(mcp_bucket, mcp_source_folder)
    source_files = []
    for file in source_folder:
        if file.endswith(".json"):
            source_files.append(file)

 mcp_ingestion_to_bq = GCSToBigQueryOperator(
        task_id="mcp_ingestion_to_bq",
        retries=0,
        dag=dag,
        bucket=mcp_bucket,
        source_objects=f"{mcp_source_folder}*.json",
        source_format="NEWLINE_DELIMITED_JSON",
        #skip_leading_rows=16,
        destination_project_dataset_table=destination_bq_table,
        write_disposition="WRITE_TRUNCATE",
        create_disposition="CREATE_NEVER",
        
        autodetect="False"
    )

4. 报错信息

File "/opt/python3.8/lib/python3.8/site-packages/google/api_core/future/polling.py", line 135, in result
    raise self._exception
google.api_core.exceptions.BadRequest: 400 Error while reading data, error message: JSON table encountered too many errors, giving up. Rows: 1; errors: 1. Please look into the errors[] collection for more details. File: gs://st-vix-ott-dev-data-usea1-akta/mcp-tr-mapping/datalatetest1.json
[2022-08-31 04:05:34,411] {taskinstance.py:1511} INFO - Marking task as FAILED. dag_id=test_mcp_ingestion, task_id=mcp_ingestion_to_bq, execution_date=20220831T040508, start_date=20220831T040530, end_date=20220831T040534
[2022-08-31 04:05:34,592] {local_task_job.py:151} INFO - Task exited with return code 1
[2022-08-31 04:05:34,687] {local_task_job.py:261} INFO - 1 downstream tasks scheduled from follow-on schedule check
排查与解决方案

可能的报错原因及对应解决方法

1. 键名与列名不匹配

JSON中包含特殊格式键名(如[CSv2] external_id、Content Type、TMS ID),这类带空格、方括号的键名,若BigQuery列名未做对应处理,会导致映射失败。

解决方法:

  • 确保BigQuery表列名与JSON键名完全一致,包含特殊字符和空格。BigQuery中含特殊字符的列名需用反引号包裹(如`[CSv2] external_id`),创建表时需正确定义。
  • 或在GCSToBigQueryOperator中添加schema_fields参数显式指定字段映射,示例:
schema_fields = [
    {"name": "ID", "type": "STRING"},
    {"name": "Title", "type": "STRING"},
    {"name": "[CSv2] external_id", "type": "STRING"},
    # 其余字段依次定义
]

mcp_ingestion_to_bq = GCSToBigQueryOperator(
    # 其他参数...
    schema_fields=schema_fields
)

2. 数据类型不匹配

例如JSON中[CSv2] pub_win_US_begin是字符串类型时间戳("1661842800"),若BigQuery对应列定义为INT64或TIMESTAMP会触发转换错误;TMS ID值为null,若列未设置允许NULL也会报错。

解决方法:

  • 逐一核对BigQuery列数据类型与JSON对应值的类型:
    • 字符串类型数字(如时间戳)可转为数值类型后导入,或直接将BigQuery列设为STRING;若需存为TIMESTAMP,可后续用SQL转换,或在JSON中存储为ISO格式时间。
    • 确保允许NULL的列在BigQuery中设置NULLABLE属性。

3. JSON格式不符合要求

NEWLINE_DELIMITED_JSON要求每行一个JSON对象,若文件存在多行JSON、语法错误(如逗号缺失、引号不闭合),会导致解析失败。

解决方法:

  • 检查GCS中JSON文件,确保每个JSON对象单独占一行,无语法错误。可使用jq工具验证:
cat gs://your-bucket/path/file.json | jq .

4. 获取详细错误信息

BigQuery报错提示提到errors[] collection,可通过以下方式查看具体错误:

  • 在Airflow任务日志中查找更详细输出;
  • 直接在BigQuery控制台手动导入该JSON文件,查看具体报错字段和原因。
验证步骤
  1. 手动在BigQuery控制台导入目标JSON文件,定位具体错误字段。
  2. 根据错误调整BigQuery表结构或JSON数据格式。
  3. 修改DAG代码,添加schema_fields确保字段映射正确。
  4. 重新运行DAG测试。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 16:09:23