如何避免Airflow中「cannot mix list...」等XCom转换错误?
Airflow 2.8.2 中XCom传递Pandas DataFrame字典报错的解决方法
问题背景
在Airflow 2.8.2(Python 3.11)的DAG中编写了如下任务函数:
@task def get_result_from_records_api(api, tries: int, result_from_salons_api: list): salon_ids_list = result_from_salons_api[1] # result is 4 pd.DataFrames result_from_records_api = get_data_or_raise_error_with_retry(api.get_records_staff_sales_of_services_goods, tries, salon_ids_list=salon_ids_list) # make lists for dfs records_lst = [] staff_from_records_lst = [] services_sales_lst = [] good_sales_lst = [] # put dfs in lists records_lst.append(result_from_records_api[0]) staff_from_records_lst.append(result_from_records_api[1]) services_sales_lst.append(result_from_records_api[2]) good_sales_lst.append(result_from_records_api[3]) # make dict with lists result_dict = { 'records': records_lst , 'staff_from_records': staff_from_records_lst , 'services_sales': services_sales_lst , 'good_sales': good_sales_lst } return result_dict
该任务可正确返回包含4个pandas.DataFrame的字典,但尝试传递结果时触发如下错误:
{xcom.py:664} ERROR - ('cannot mix list and non-list, non-null values', 'Conversion failed for column position with type object'). If you are using pickle instead of JSON for XCom, then you need to enable pickle support for XCom in your airflow config or make sure to decorate your object with attr.
已在docker-compose.yml中设置AIRFLOW__CORE__ENABLE_XCOM_PICKLING=true,寻求解决方法。
解决方法
1. 确认Airflow配置是否生效
- 重启所有Airflow服务,确保
docker-compose.yml中的配置变更已加载 - 进入Airflow WebUI的Admin > Configurations页面,验证
core.enable_xcom_pickling参数是否显示为True
2. 简化返回结构,移除不必要的嵌套列表
当前代码将每个DataFrame放入单独列表后再存入字典,这种嵌套结构可能干扰Pickle序列化逻辑。直接将DataFrame作为字典的值即可:
@task def get_result_from_records_api(api, tries: int, result_from_salons_api: list): salon_ids_list = result_from_salons_api[1] result_from_records_api = get_data_or_raise_error_with_retry(api.get_records_staff_sales_of_services_goods, tries, salon_ids_list=salon_ids_list) # 直接用DataFrame作为字典值,无需嵌套列表 result_dict = { 'records': result_from_records_api[0], 'staff_from_records': result_from_records_api[1], 'services_sales': result_from_records_api[2], 'good_sales': result_from_records_api[3] } return result_dict
3. 转换DataFrame为可序列化的字符串格式
如果Pickle仍有问题,可以将DataFrame转换为CSV字符串存储,下游任务再解析回DataFrame:
# 任务函数中转换为CSV字符串 from io import StringIO @task def get_result_from_records_api(api, tries: int, result_from_salons_api: list): salon_ids_list = result_from_salons_api[1] result_from_records_api = get_data_or_raise_error_with_retry(api.get_records_staff_sales_of_services_goods, tries, salon_ids_list=salon_ids_list) result_dict = { 'records': result_from_records_api[0].to_csv(index=False), 'staff_from_records': result_from_records_api[1].to_csv(index=False), 'services_sales': result_from_records_api[2].to_csv(index=False), 'good_sales': result_from_records_api[3].to_csv(index=False) } return result_dict # 下游任务中解析CSV字符串恢复DataFrame def downstream_task(ti): import pandas as pd xcom_data = ti.xcom_pull(task_ids='get_result_from_records_api') records_df = pd.read_csv(StringIO(xcom_data['records']))
4. 清洗DataFrame中的异常列数据
错误提示提到Conversion failed for column position with type object,说明DataFrame的position列存在混合列表与非列表值的异常数据:
- 检查
result_from_records_api返回的DataFrame,查看position列的内容,确认是否有列表类型的单元格 - 对该列进行清洗,例如将列表转换为字符串:
df['position'] = df['position'].apply(lambda x: str(x) if isinstance(x, list) else x)
内容的提问来源于stack exchange,提问作者John Doe
相关产品推荐
相关产品推荐

