如何通过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属性。
- 字符串类型数字(如时间戳)可转为数值类型后导入,或直接将BigQuery列设为
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文件,查看具体报错字段和原因。
验证步骤
- 手动在BigQuery控制台导入目标JSON文件,定位具体错误字段。
- 根据错误调整BigQuery表结构或JSON数据格式。
- 修改DAG代码,添加
schema_fields确保字段映射正确。 - 重新运行DAG测试。
内容的提问来源于stack exchange,提问作者Priya
相关产品推荐
相关产品推荐

