You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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:

  1. 创建自定义序列化器文件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)
  1. 在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 17:50:29