This commit is contained in:
doudou201703
2026-08-02 14:07:04 +08:00
commit 046fac6dc1
339 changed files with 19435 additions and 0 deletions
View File
+38
View File
@@ -0,0 +1,38 @@
"""配置文件"""
# broker_url = "redis://default:com792282@127.0.0.1:6379/5"
# result_backend = "redis://:com792282@127.0.0.1:6379/6"
# result_backend = "django-db" # 任务结果存储数据库:这里默认是存储到数据库中
CELERY_BROKER_CONNECTION_RETRY_ON_STARTUP = True
# 用于设置单个任务的最大允许执行时间(单位通常是秒)。当 Celery Worker 执行的任务耗时超过这个设定值时,Celery将会强制终止该任务的执行。这个机制可以防止某个任务无限期地运行下去,从而可能导致Worker阻塞无法处理新的任务。
CELERY_TASK_TIME_LIMIT = 30 * 60
timezone = "Asia/Shanghai" # 时区:与 django 使用相同的时区;
# 这个配置项设置了Celery在发布任务和存储任务结果时使用的序列化器。当Celery Worker从Broker接收任务并将任务结果返回到Broker时,都会用到这个序列化器。在您给出的配置中,CELERY_TASK_SERIALIZER = "json" 表明任务将以JSON格式进行序列化。JSON是一种轻量级的数据交换格式,易于阅读和编写,同时也便于不同语言之间的解析。
task_serializer = "json"
# 此配置项定义了Celery Worker所接受的任务消息内容类型(MIME类型)。换句话说,只有在这个列表中的内容类型的任务消息,Celery Worker才会接受并尝试去执行。在您给出的配置中,CELERY_ACCEPT_CONTENT = ['json']
# 表示Worker仅接受以JSON格式序列化的任务消息。这与CELERY_TASK_SERIALIZER配置相呼应,确保系统内部的一致性,即发布的任务消息和接收的任务消息都采用相同的序列化格式。
accept_content = [ "json" ]
# CELERY_BEAT_SCHEDULER = "django_celery_beat.schedulers:DatabaseScheduler"
# 下面的语法是实现 定时作务的; USE_TZ
# DJANGO_CELERY_BEAT_TZ_AWARE = False # 一定要配置,否则会报mysql不支持的时间设置
# CELERY_BEAT_SCHEDULER = "django_celery_beat.schedulers:DatabaseScheduler" # 默认配置
# CELERY_BEAT_SCHEDULE = {
# "add_demo_event_one_hours": {
# "task": "django_celery.scheduledUpdatestoken.tasks.updateToken",
# "schedule": timedelta( seconds=60 )
# },
#
# }
"""
django_celery_beat.models.ClockedSchedule # 此模型存放已经关闭的任务
django_celery_beat.models.CrontabSchedule # cron的时间表
django_celery_beat.models.IntervalSchedule # 以特定间隔(例如,每5秒)运行的计划。以特定时间间隔(例如每 5 秒)运行的计划。
django_celery_beat.models.PeriodicTask # 此模型定义要运行的单个周期性任务。
django_celery_beat.models.PeriodicTasks # 此模型仅用作索引以跟踪计划何时更改
django_celery_beat.models.SolarSchedule # 定制任务
from django_celery_beat.models import PeriodicTask, CrontabSchedule 包含 cron 中的条目等字段的计划:分钟、小时、星期几、day_of_month month_of_year。
import json
"""
+152
View File
@@ -0,0 +1,152 @@
"""
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
+60
View File
@@ -0,0 +1,60 @@
"""注释"""
import requests
from celery import Task
from channels.layers import get_channel_layer
from django_celery.main import app
from wechat_token.token工具包 import sixToken
from wechat_数据库.支付订单数据库.商品订单数据 import order订单
from wechat_项目工具包.小程序交易管理服务.小程序交易订单管理工具 import miniTransaction交易管理服务工具
channel_layer = get_channel_layer()
import logging
logger = logging.getLogger(__name__)
# 这是任务执行成功后的回调监听函数;(使用方法: 只需要在调用任务的时候添加属性 link=函数名.s() )
@app.task(bind=True)
def on_success_callback(self: Task, result, *args, **kwargs):
"""这是任务执行成功后的回调监听函数"""
# 在这里实现你的回调逻辑
print("成功地回调")
task_name = self.name
# 获取触发回调的任务ID
task_id = self.request.id
print(self.request.parent_id) # 这个参数非常重要,是执行任务的🆔(可以知道是哪个任务的回调)
print(f"任务 {task_name} (ID: {task_id}) 执行成功,回调函数已触发")
class MyCustomTask(Task):
"""注释"""
def on_success(self, retval, task_id, args, kwargs):
""" 任务执行成功,回调函数打印返回值 """
# 在任务成功执行后运行的逻辑
# logger.info( f"任务 {task_id} 成功: {retval}" )
pass
def on_retry(self, exc, task_id, args, kwargs, einfo):
""" 在任务重试后运行的逻辑 """
pass
def on_failure(self, exc, task_id, args, kwargs, einfo):
"""参数"""
# 在任务执行失败后运行的逻辑
pass
@app.task(base=MyCustomTask, bind=True)
def mini调用微信小程序交易管理服务录入订单接口(self: Task,order支付商家侧订单号, **kwargs, ):
"""执行这个任务,:会自动打印回调"""
# 任务完成后,往前端发消息 print("任务执行了")
# TODO my_task.apply_async( kwargs={ "name": "这里是传递的名称的值" }, queue="xiao_queue", )
# 根据订单号: 对运单接口发起生成运单的物流信息
print("\n🍎 🍏 🍊 🍎 🍏 🍊 🍎 🍏 🍊 到底调用了没有呀", )
miniTransaction交易管理服务工具().发货信息录入(order支付商家侧订单号=order支付商家侧订单号)
+61
View File
@@ -0,0 +1,61 @@
"""注释"""
from celery import Task
from django_celery.main import app
import logging
from wechat_数据库.雪花算法.雪花算法工具 import generator
logger = logging.getLogger(__name__)
# 这是任务执行成功后的回调监听函数;(使用方法: 只需要在调用任务的时候添加属性 link=函数名.s() )
@app.task(bind=True)
def on_success_callback(self: Task, result, *args, **kwargs):
"""这是任务执行成功后的回调监听函数"""
# 在这里实现你的回调逻辑
print("成功地回调")
task_name = self.name
# 获取触发回调的任务ID
task_id = self.request.id
print(self.request.parent_id) # 这个参数非常重要,是执行任务的🆔(可以知道是哪个任务的回调)
print(f"任务 {task_name} (ID: {task_id}) 执行成功,回调函数已触发")
class MyCustomTaskMessages(Task):
"""注释"""
def on_success(self, retval, task_id, args, kwargs):
""""""
# 在任务成功执行后运行的逻辑
# logger.info( f"任务 {task_id} 成功: {retval}" )
pass
# print( retval )
def on_retry(self, exc, task_id, args, kwargs, einfo):
""""""
# 在任务重试后运行的逻辑
pass
def on_failure(self, exc, task_id, args, kwargs, einfo):
""""""
pass
# 只记录,不要 raise 非法对象!
print(f"Task failed: {exc}")
# 不要在这里 raise 任何东西,除非是合法异常
@app.task(base=MyCustomTaskMessages, bind=True, queue="xiao_queue", retry_kwargs={'max_retries': 5}, retry_backoff=True, )
def taskMessages(self, xiaoziziname, **kwargs):
""""""
try:
pass
auto_str = generator.generate_str()
print("成功了", auto_str, xiaoziziname)
return ""
except Exception as e:
pass
print(e)
return e
+32
View File
@@ -0,0 +1,32 @@
from django.core.handlers.wsgi import WSGIRequest
from rest_framework import status
from rest_framework.decorators import action
from rest_framework.request import Request
from rest_framework.response import Response
from rest_framework.routers import DefaultRouter
from rest_framework.viewsets import ViewSet
from django_celery.xiao_queue队列.tasks import taskMessages
from sixGuoDjango import BaseResponse, ResData
class Celery测试接口ViewSets( ViewSet ):
@action( detail=False, methods=[ "get", "post", "put", "delete" ], url_path="test" )
def 测试路由接口( self, request: WSGIRequest | Request ):
""""""
taskMessages.delay( token="alkdsfjal;d" )
taskMessages.apply_async( kwargs={ 'token': 'your_token' }, )
httpResData = BaseResponse( resData=ResData( listData=[ "" ], dictData={ "demo": "xiaozizi" } ), ).model_dump( )
return Response( data=httpResData, status=status.HTTP_200_OK )
Celery测试接口ViewSetsRouter = DefaultRouter( )
Celery测试接口ViewSetsRouter.register( prefix="", viewset=Celery测试接口ViewSets, basename="" )