153 lines
5.1 KiB
Python
153 lines
5.1 KiB
Python
"""
|
||
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
|