使用Airflow Cloud Composer加载Parquet至BigQuery:单列未加载问题
解决Parquet文件追加至BigQuery时某列未加载的问题
针对你遇到的问题,按以下步骤排查和解决:
1. 移除无效参数skip_leading_rows=1
skip_leading_rows仅适用于CSV格式文件,对Parquet列式存储文件完全无效,保留该参数可能干扰数据读取逻辑。修改Operator配置,删除这一行:
GCS_to_BQ = GoogleCloudStorageToBigQueryOperator( task_id='Loading_Parquet_Into_BQ', bucket='my bucket', source_objects=['path/to/parquet/year=2021/month=8/*.parquet', 'path/to/parquet/year=2022/month=9/*.parquet'], source_format='parquet', destination_project_dataset_table='project_name.dataset_name.table_name', create_disposition='CREATE_IF_NEEDED', write_disposition='WRITE_APPEND', schema_update_options=['ALLOW_FIELD_RELAXATION', 'ALLOW_FIELD_ADDITION'], autodetect=True, # 移除下面这行无效参数 # skip_leading_rows=1, google_cloud_storage_conn_id='google_cloud_default', bigquery_conn_id='google_cloud_default' )
2. 验证Parquet文件中的目标列
- 用BigQuery命令行工具检查Parquet文件元数据,确认目标列存在:
bq show --format=prettyjson gs://my-bucket/path/to/parquet/your-target-file.parquet - 确认该列在所有待加载的Parquet文件中都存在,且包含非NULL数据(全NULL的列可能无法被正确识别或加载)。
3. 检查列名匹配问题
BigQuery对列名大小写不敏感,但存储时保留原始大小写。确保Parquet文件中的列名与BigQuery表中的列名完全一致(包括大小写),避免因大小写差异导致映射失败。
4. 手动指定Schema替代自动检测
如果自动检测仍无法识别目标列,尝试手动指定Schema:
- 导出BigQuery表当前的Schema,补充缺失列的定义
- 在Operator中配置
schema_fields参数,关闭autodetect:schema_fields = [ {"name": "col1", "type": "STRING", "mode": "NULLABLE"}, {"name": "col2", "type": "INTEGER", "mode": "NULLABLE"}, {"name": "missing_col", "type": "FLOAT", "mode": "NULLABLE"}, # 替换为你的列信息 # 其他列... ] GCS_to_BQ = GoogleCloudStorageToBigQueryOperator( task_id='Loading_Parquet_Into_BQ', bucket='my bucket', source_objects=['path/to/parquet/year=2021/month=8/*.parquet', 'path/to/parquet/year=2022/month=9/*.parquet'], source_format='parquet', destination_project_dataset_table='project_name.dataset_name.table_name', create_disposition='CREATE_IF_NEEDED', write_disposition='WRITE_APPEND', schema_update_options=['ALLOW_FIELD_RELAXATION', 'ALLOW_FIELD_ADDITION'], # 关闭自动检测 # autodetect=True, schema_fields=schema_fields, google_cloud_storage_conn_id='google_cloud_default', bigquery_conn_id='google_cloud_default' )
5. 查看任务日志排查细节
在Airflow任务日志中搜索目标列的相关信息,检查是否存在数据类型不匹配、列映射失败等警告或错误,根据日志进一步定位问题。
内容的提问来源于stack exchange,提问作者Moushmi Das
相关产品推荐
相关产品推荐

