Polars DataFrame写入BigQuery时Struct列模式不匹配报错排查
问题描述
将基于Pandas的代码迁移为使用Polars DataFrame后,无法将数据写入BigQuery表,触发Schema不匹配错误。使用Pandas时数据可正常写入目标表。
错误信息
google.api_core.exceptions.BadRequest: 400 Error while reading data,
error message:
Schema mismatch: referenced variable 'items.list.item.id_item' has array levels of 1,
while the corresponding field path to Parquet column has 0 repeated fields;
reason: invalid,
message: Error while reading data,
error message:
Schema mismatch: referenced variable 'items.list.item.id_item' has array levels of 1,
while the corresponding field path to Parquet column has 0 repeated fields
代码示例
import io from typing import Dict, List, Union import pandas as pd import polars as pl from google.cloud import bigquery def append_dataframe_to_table( data: Union[pd.DataFrame, pl.DataFrame], bq_client: bigquery.Client, table_schema: List[Dict[str, str]], destination_table: str, create_disposition: str, write_disposition: str, wait_until_finished: bool, ) -> bigquery.LoadJob: """ Append a DataFrame to a BigQuery table. :param data: Data to be stored. :param table_schema: Intended schema for the target table. Example: table_schema = [ {'name': 'id', 'type': 'INT64', 'mode': 'REQUIRED'}, {'name': 'full_name', 'type': 'STRING', 'mode': 'REQUIRED'}, {'name': 'date_of_birth', 'type': 'DATETIME', 'mode': 'NULLABLE'} ] :param destination_table: Full path of the destination table. :param create_disposition: Describes conditions when a job should create a table. Possible values: |--------------------------------|-------------------------------------------------------------| | CREATE_DISPOSITION_UNSPECIFIED | Unknown. | | CREATE_NEVER | This job should never create tables. | | CREATE_IF_NEEDED | This job should create a table if it doesn't already exist. | :param write_disposition: Describes whether a mutation to a table should overwrite or append. Possible values: |-------------------------------|-----------------------------------------------------------------| | WRITE_DISPOSITION_UNSPECIFIED | Unknown. | | WRITE_EMPTY | This job should only be writing to empty tables. | | WRITE_TRUNCATE | This job will truncate table data and write from the beginning. | | WRITE_APPEND | This job will append to a table. | :param wait_until_finished: If true, wait for the job to finish. """ job_config = bigquery.LoadJobConfig() job_config.create_disposition = create_disposition job_config.write_disposition = write_disposition job_config.schema = table_schema if isinstance(data, pl.DataFrame): job_config.source_format = bigquery.SourceFormat.PARQUET # Write DataFrame to stream as parquet file; does not hit disk with io.BytesIO() as stream: data.write_parquet(stream) stream.seek(0) load_job = bq_client.load_table_from_file( file_obj=stream, destination=destination_table, job_config=job_config ) else: load_job = bq_client.load_table_from_dataframe( dataframe=data, destination=destination_table, job_config=job_config ) if wait_until_finished: print("Waiting for load job to finish.") load_job.result() return load_job if __name__ == "__main__": table_schema = [ {"name": "hash", "type": "STRING", "mode": "REQUIRED"}, { "name": "items", "type": "RECORD", "mode": "REPEATED", "fields": [ {"name": "id_item", "type": "INT64", "mode": "NULLABLE"}, {"name": "id_category", "type": "INT64", "mode": "NULLABLE"}, {"name": "manufacturer", "type": "STRING", "mode": "NULLABLE"}, ], }, {"name": "prediction", "type": "BOOLEAN", "mode": "NULLABLE"}, {"name": "confidence", "type": "FLOAT64", "mode": "NULLABLE"}, ] df = pl.DataFrame( { "hash": ["abcd", "efg"], "items": [ [ {"id_item": 1234, "id_category": 12, "manufacturer": "A"}, {"id_item": 1235, "id_category": 12, "manufacturer": "B"}, ], [{"id_item": 2345, "id_category": 13, "manufacturer": "C"}], ], "prediction": [True, False], "confidence": [0.4, None], } ) bq_client = bigquery.Client(project="my-project") append_dataframe_to_table( data=df, bq_client=bq_client, table_schema=table_schema, destination_table="my-project.my_dataset.my_test_table", create_disposition="CREATE_IF_NEEDED", write_disposition="WRITE_APPEND", wait_until_finished=True, )
解决方案
问题根源
Polars在处理**嵌套数组(对应BigQuery的REPEATED RECORD类型)**时,默认写入Parquet的格式与Pandas存在差异。BigQuery期望items是数组包裹的结构体,但当前Polars生成的Parquet中,items的嵌套结构未被正确识别为数组内的对象,导致Schema不匹配。
方法1:显式定义Polars嵌套Schema(推荐)
在创建Polars DataFrame时,明确指定items列的类型为结构体数组,确保写入Parquet时保留正确的嵌套层级:
修改__main__中创建DataFrame的代码:
# 显式定义items列的结构体Schema items_struct = pl.Struct([ pl.Field("id_item", pl.Int64), pl.Field("id_category", pl.Int64), pl.Field("manufacturer", pl.Utf8) ]) df = pl.DataFrame( { "hash": ["abcd", "efg"], "items": pl.Series( [ [ {"id_item": 1234, "id_category": 12, "manufacturer": "A"}, {"id_item": 1235, "id_category": 12, "manufacturer": "B"}, ], [{"id_item": 2345, "id_category": 13, "manufacturer": "C"}], ], dtype=pl.Array(items_struct) # 指定为结构体数组 ), "prediction": [True, False], "confidence": [0.4, None], } )
方法2:转换为Pandas DataFrame后写入
如果不想调整Polars数据结构,可先将Polars DataFrame转为Pandas DataFrame,复用Pandas与BigQuery兼容的嵌套结构处理逻辑:
修改append_dataframe_to_table函数中的Polars分支:
if isinstance(data, pl.DataFrame): # 转换为Pandas DataFrame后使用原生方法写入 load_job = bq_client.load_table_from_dataframe( dataframe=data.to_pandas(), destination=destination_table, job_config=job_config ) else: load_job = bq_client.load_table_from_dataframe( dataframe=data, destination=destination_table, job_config=job_config )
注意:此方法适合数据量较小的场景,避免转换带来的性能开销。
额外注意事项
- 确保BigQuery目标表的Schema中,
items字段的mode为REPEATED、type为RECORD,内部字段定义与数据完全匹配。 - 方法1是Polars原生方案,能更好保留数据类型精度,避免跨库转换的潜在问题。
内容的提问来源于stack exchange,提问作者Michal Chromčák

