Airflow读取BigQuery表转DataFrame出现_pickle.PicklingError如何解决
问题根因
- 代码中存在重复调用
bqclient.query()的逻辑错误:第一次调用该方法已经返回QueryJob对象,你又将这个对象作为参数传入了第二次bqclient.query()调用,导致BigQuery客户端实例被带入配置深拷贝流程。Google云客户端本身不支持pickle序列化,Airflow运行时的序列化机制触发了这个限制,因此抛出对应报错。
修复方案
1 修正重复调用的逻辑错误
直接修改读取逻辑,仅调用一次bqclient.query()即可,修复后代码如下:
import pandas as pd from google.cloud import bigquery as bq from google.oauth2 import service_account def reading(ds, **kwargs): project_id = '{project-id}' credentials = service_account.Credentials.from_service_account_file('/usr/local/airflow/extras/credentials.json') bqclient = bq.Client(credentials=credentials, project=project_id) # 仅定义SQL字符串,不提前执行query query_sql = """ SELECT * FROM `{project-id}.{schema-name}.{table-name}` """ # 单次调用query后直接转DataFrame dataframe = ( bqclient.query(query_sql) .result() .to_dataframe( create_bqstorage_client=True, ) ) print(dataframe.head()) # 不需要返回值的话显式返回None,避免不必要的序列化 return
2 其他优化建议
- 如果不需要将DataFrame传递给后续任务,给调用该函数的PythonOperator设置
do_xcom_push=False参数,彻底关闭该任务的XCom序列化逻辑,降低序列化冲突概率 - 如果需要将DataFrame传递给后续任务,不要直接返回DataFrame对象,可将其写入CSV/Parquet格式文件存储到共享存储、GCS等位置,后续任务读取文件即可
- 推荐直接使用Airflow官方提供的
BigQueryHook完成读取操作,Hook已经适配了Airflow的序列化规则,不需要自行处理客户端初始化逻辑,示例代码如下:
from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook def reading(ds, **kwargs): # 提前在Airflow连接管理中配置好GCP连接,填对应连接ID即可 bq_hook = BigQueryHook(gcp_conn_id='your_gcp_conn_id', use_legacy_sql=False) dataframe = bq_hook.get_pandas_df(sql=""" SELECT * FROM `{project-id}.{schema-name}.{table-name}` """) print(dataframe.head()) return
内容的提问来源于stack exchange,提问作者Saurav Ganguli
相关产品推荐
相关产品推荐

