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

Celery模式Airflow异步运行自定义任务遇NotRegistered错误求助

解决Airflow Celery Worker无法识别自定义Celery任务的问题

问题背景

部署了基于Celery+Redis的Airflow环境,默认DAG任务由Celery Worker执行。尝试在某个DAG任务中调用自定义Celery任务时,触发celery.exceptions.NotRegistered: 'maximum'错误。

相关代码与日志

自定义任务代码(dags/tasks.py)

from airflow.configuration import conf
from airflow.config_templates.default_celery import DEFAULT_CELERY_CONFIG
from celery import Celery
from celery import shared_task

if conf.has_option('celery', 'celery_config_options'):
    celery_configuration = conf.getimport('celery', 'celery_config_options')
else:
    celery_configuration = DEFAULT_CELERY_CONFIG

app = Celery(conf.get('celery', 'CELERY_APP_NAME'), config_source=celery_configuration,include=["dags.tasks"])
app.autodiscover_tasks(force=True)
print("here")
print(conf.get('celery', 'CELERY_APP_NAME'))
print(celery_configuration)
print(app)
@app.task(name='maximum')
def maximum(x=10, y=11):
    print(x)
    if x > y:
        return x
    else:
        return y

tasks = app.tasks.keys()
print(tasks)

DAG中调用任务的代码

max=maximum.apply_async(kwargs={'x':5, 'y':4})
print(max)
print(max.get(timeout=5))

错误信息

File "/home/airflow/.local/lib/python3.7/site-packages/celery/result.py", line 336, in maybe_throw
    self.throw(value, self._to_remote_traceback(tb))
File "/home/airflow/.local/lib/python3.7/site-packages/celery/result.py", line 329, in throw
    self.on_ready.throw(*args, **kwargs)
File "/home/airflow/.local/lib/python3.7/site-packages/vine/promises.py", line 234, in throw
    reraise(type(exc), exc, tb)
File "/home/airflow/.local/lib/python3.7/site-packages/vine/utils.py", line 30, in reraise
    raise value
celery.exceptions.NotRegistered: 'maximum'

关键现象

  • 本地查询已注册任务时,能看到maximum:
    dict_keys(['celery.chunks', 'airflow.executors.celery_executor.execute_command', 'maximum', 'celery.backend_cleanup', 'celery.chord_unlock', 'celery.group', 'celery.map', 'celery.accumulate', 'celery.chain', 'celery.starmap', 'celery.chord'])
    
  • Airflow Worker日志显示仅注册了airflow.executors.celery_executor.execute_command任务,无自定义的maximum任务。

问题根源

Airflow的Celery Worker默认加载的是Airflow自带的Celery App(airflow.executors.celery_executor.app),而非自定义的Celery App。Worker未加载tasks.py模块,因此无法识别自定义任务。

解决方案

方案1:通过Airflow插件加载自定义任务

  1. 在Airflow的plugins目录下创建插件文件(如custom_celery_tasks.py):
    from airflow.plugins_manager import AirflowPlugin
    from celery import shared_task
    
    @shared_task(name='maximum')
    def maximum(x=10, y=11):
        print(x)
        return x if x > y else y
    
    class CustomCeleryTasksPlugin(AirflowPlugin):
        name = "custom_celery_tasks"
    
  2. 重启Airflow Worker和Scheduler,让插件被加载:
    docker-compose restart airflow-worker airflow-scheduler
    
  3. 在DAG中直接调用maximum.apply_async即可,任务会自动注册到Airflow的Celery App中。

方案2:修改Worker启动命令加载自定义任务

  1. 修改docker-compose.yaml中airflow-worker的启动命令,添加--include参数指定任务模块:
    airflow-worker:
        <<: *airflow-common
        command: celery worker --app airflow.executors.celery_executor.app --include dags.tasks
        healthcheck:
          test:
            - "CMD-SHELL"
            - 'celery --app airflow.executors.celery_executor.app inspect ping -d "celery@$${HOSTNAME}"'
          interval: 10s
          timeout: 10s
          retries: 5
        restart: always
    
  2. 重启Worker:
    docker-compose restart airflow-worker
    
  3. 重启后Worker会加载dags.tasks模块,自定义任务maximum会被注册。

方案3:使用Airflow PythonOperator封装任务(推荐)

若无需直接调用Celery任务,可直接用Airflow的PythonOperator封装函数,完全适配Airflow执行模式:

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime

def maximum(x=10, y=11):
    print(x)
    return x if x > y else y

with DAG('custom_task_dag', start_date=datetime(2022,9,2)) as dag:
    maximum_task = PythonOperator(
        task_id='calculate_max',
        python_callable=maximum,
        op_kwargs={'x':5, 'y':4}
    )

内容的提问来源于stack exchange,提问作者atul.mishra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 06:03:18