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插件加载自定义任务
- 在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" - 重启Airflow Worker和Scheduler,让插件被加载:
docker-compose restart airflow-worker airflow-scheduler - 在DAG中直接调用
maximum.apply_async即可,任务会自动注册到Airflow的Celery App中。
方案2:修改Worker启动命令加载自定义任务
- 修改
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 - 重启Worker:
docker-compose restart airflow-worker - 重启后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
相关产品推荐
相关产品推荐

