Airflow Docker环境运行含Numpy的DAG失败问题求助
解决方案
一、Docker环境下配置Airflow使用Pickle序列化(无需手动生成airflow.cfg)
Airflow Docker部署时,默认不会在容器内生成可编辑的airflow.cfg,推荐通过环境变量或挂载自定义配置文件修改序列化设置:
1. 环境变量快速配置
在你的docker-compose.yml中给webserver和scheduler服务添加以下环境变量:
version: '3.8' services: airflow-webserver: environment: - AIRFLOW__CORE__ENABLE_XCOM_PICKLING=True - AIRFLOW__CORE__XCOM_SERIALIZER=pickle - AIRFLOW__CORE__XCOM_DESERIALIZER=pickle airflow-scheduler: environment: - AIRFLOW__CORE__ENABLE_XCOM_PICKLING=True - AIRFLOW__CORE__XCOM_SERIALIZER=pickle - AIRFLOW__CORE__XCOM_DESERIALIZER=pickle
修改后重启容器:docker-compose up -d --force-recreate
2. 挂载自定义airflow.cfg
若需更复杂配置,本地创建airflow.cfg文件并添加以下内容:
[core] enable_xcom_pickling = True xcom_serializer = pickle xcom_deserializer = pickle
在docker-compose.yml中挂载该文件到容器内对应路径:
services: airflow-webserver: volumes: - ./airflow.cfg:/opt/airflow/airflow.cfg airflow-scheduler: volumes: - ./airflow.cfg:/opt/airflow/airflow.cfg
重启容器即可生效。
二、安全替代方案:避免直接传递numpy ndarray
Pickle序列化存在安全风险(反序列化可能执行恶意代码),推荐以下两种更安全的方式:
1. 转换ndarray为原生列表传递
在DAG任务中将numpy数组转为Python列表,完成XCom传递后再还原:
import numpy as np from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def generate_array(): arr = np.array([1, 2, 3, 4]) return arr.tolist() # 转为可JSON序列化的列表 def process_array(**context): arr_list = context['ti'].xcom_pull(task_ids='generate_array') arr = np.array(arr_list) # 还原为ndarray print(f"Processed array: {arr}") with DAG('numpy_xcom_demo', start_date=datetime(2024, 1, 1), schedule_interval=None) as dag: t1 = PythonOperator(task_id='generate_array', python_callable=generate_array) t2 = PythonOperator(task_id='process_array', python_callable=process_array, provide_context=True) t1 >> t2
2. 自定义XCom序列化器
编写专属序列化逻辑,仅处理numpy类型,无需全局启用Pickle:
- 创建自定义序列化器文件
custom_xcom_serializers.py:
import numpy as np from airflow.serialization.serializers.base_serializer import BaseSerializer class NumpyArraySerializer(BaseSerializer): @classmethod def serialize(cls, value): if isinstance(value, np.ndarray): return { '__type': 'numpy.ndarray', 'data': value.tolist(), 'dtype': str(value.dtype) } return super().serialize(value) @classmethod def deserialize(cls, value): if value.get('__type') == 'numpy.ndarray': return np.array(value['data'], dtype=value['dtype']) return super().deserialize(value)
- 在docker-compose.yml中挂载文件并配置环境变量:
services: airflow-webserver: volumes: - ./custom_xcom_serializers.py:/opt/airflow/custom_xcom_serializers.py environment: - AIRFLOW__CORE__XCOM_SERIALIZERS=custom_xcom_serializers.NumpyArraySerializer,airflow.serialization.serializers.json_serializer.JsonSerializer airflow-scheduler: volumes: - ./custom_xcom_serializers.py:/opt/airflow/custom_xcom_serializers.py environment: - AIRFLOW__CORE__XCOM_SERIALIZERS=custom_xcom_serializers.NumpyArraySerializer,airflow.serialization.serializers.json_serializer.JsonSerializer
重启容器后即可直接传递ndarray对象。
内容的提问来源于stack exchange,提问作者sogu
相关产品推荐
相关产品推荐

