""" Celery 启动入口文件 - 必须放在最顶层 """ import os import django from sixGuoDjango.settings import BASE_DIR os.environ.setdefault( key="DJANGO_SETTINGS_MODULE", value=f"{BASE_DIR.name}.settings" ) django.setup( ) import orjson from kombu import Queue, Exchange from kombu.serialization import register from celery import Celery, signals # redis_url = os.environ.get("REDIS_URL") redis_url = "127.0.0.1:6379" # redis_url = f"{REDIS_IP}:6379" # 作用就是:需要加载一个django项目的配置文件# 为"芹菜"程序设置默认的Django设置模块。; # app.conf.beat_schedule = { app = Celery( main="django_celery", broker=f"redis://default:com792282@{redis_url}/9", backend=f"redis://default:com792282@{redis_url}/8", # 存储数据库 namespace="CELERY", # 作用:配置前缀(如 "CELERY"),用于从 settings 中提取配置。Django 中默认为 "CELERY" task_queues=[ Queue( name="xiao_queue", routing_key="xiao_queue", exchange=Exchange( name="xiao_queue", type="direct" ) ), Queue( name="celery", routing_key="celery", exchange=Exchange( name="celery", type="direct" ) ), # 添加默认队列 ], # 作用:预定义队列列表(kombu.Queue对象) worker_cancel_long_running_tasks_on_connection_loss=True, ## 默认,保持任务完整性 # task_acks_late=True, # 任务确认延迟 task_time_limit=60, # 任务硬性时间限制(以秒为单位)。超过此限制后,正在处理该任务的 worker 将被终止并由新的 worker 替换 # 关键配置:指定xiao_queue队列只使用一个worker+可选:确保每个任务由独立worker处理 worker_max_tasks_per_child=100, # 每个子进程执行多少任务后自动重启,防止内存泄漏。 worker_prefetch_multiplier=1, ## 每个 worker 一次只取 1 个任务(避免长任务阻塞) ) app.conf.update( timezone="Asia/Shanghai", enable_utc=False, task_serializer="utf8json", result_serializer="utf8json", accept_content=[ "json", "utf8json" ], ) json_dumps = lambda obj: orjson.dumps( obj ).decode( ) json_loads = lambda obj: orjson.loads( obj ) # 注册自定义 JSON 序列化器 register( name="utf8json", encoder=json_dumps, decoder=json_loads, content_type="application/json", content_encoding="utf-8" ) app.autodiscover_tasks( packages=[ "django_celery.message", "django_celery.xiao_queue队列", ], ) # import logging logger = logging.getLogger( __name__ ) # @task_prerun.connect # def log_task_queue( sender=None, task_id=None, task=None, args=None, kwargs=None, **kwds ): # """ # :param sender: # :type sender: # :param task_id: # :type task_id: # :param task: # :type task: # :param args: # :type args: # :param kwargs: # :type kwargs: # :param kwds: # :type kwds: # """ # # 正确方式:从 request.delivery_info 获取实际消费的队列(即 routing_key) # delivery_info = current_task.request.delivery_info # queue_name = delivery_info.get( "routing_key", "unknown" ) if delivery_info else "unknown" # logger.warning( f"[QUEUE DEBUG] Task {task.name}[{task_id}] is running from queue: {queue_name}" ) # """ # Worker1:只处理 xiao_queue,使用 1 个进程 celery -A django_celery.main worker -Q xiao_queue --concurrency=10 -n worker_xiao # Worker2:处理其他队列,使用 9 个进程 celery -A django_celery.main worker -Q celery --concurrency=1 -n worker_default worker_xiao:指定的线程名称, -Q 指定的队列的名称(--queue) -A指定app(--app) --concurrency=1(指定的数量进程或者线程) -P threads(指定使用线程) -n(--hostname)给进程或者线程取个指定的名称 celery -A django_celery.main worker -Q xiao_queue --concurrency=100 -P threads -n worker_xiao -E (推荐:最合适的启动线程) celery -A django_celery.main worker -Q xiao_queue --concurrency=1 -P prefork -n worker_xiao (不推荐配置不够强大进程) celery -A django_celery.main worker -Q sql_queue --concurrency=1 -P prefork -n worker_sql (不推荐配置不够强大进程) celery -A django_celery.main worker -Q xiao_queue --concurrency=3000 -P eventlet -n worker_xiao (这里启动的是协程) celery -A django_celery.main worker -Q xiao_queue --concurrency=3000 -P gevent -n worker_xiao (不推荐,协程,需要猴子协议) celery --app django_celery.main worker --pool=threads --concurrency=10 """ # ========== 使用Celery信号系统 ========== @signals.worker_shutdown.connect def on_worker_shutdown( sender=None, **kwargs ): """Celery worker关闭时的回调""" print( "\n🚀 Celery Worker正在优雅关闭..." ) print( "1. 保存任务状态..." ) print( "2. 关闭数据库连接..." ) print( "3. 清理临时文件..." ) @signals.worker_ready.connect def on_worker_ready( sender=None, **kwargs ): """Worker启动完成时的回调""" # 打印所有已注册任务 print( "=== 已注册任务 ===" ) for name in sender.app.tasks.keys( ): if not name.startswith( 'celery.' ): print( f" {name}" ) if __name__ == '__main__': pass