Windows Celery进阶之路:自定义基类、进度监控、周期任务与Django全整合
2026/9/13 16:45:33 网站建设 项目流程

一、celery 介绍

Celery是一个简单的、快速的,灵活且可靠的分布式系统,用于处理大量消息,同时提供了一些工具来维护这样的一个系统。这是一个专注于实时处理的任务队列,同时也支持任务调度。

Celery 支持自主配置消息队列,结果存储,并发,序列化等,这里我们使用在windows下使用redis作为消息队列和结果存储,使用eventlet并发, 序列化使用json的方式。

二、celery基本应用

  1. 安装 Python 及相关组件
pip install celery redis==7.1.3eventlet flower
  1. 新建一个主应用

文件名我们命名为 main.py

importtimefromceleryimportCelery broker_url='redis://127.0.0.1:6379/1'result_backend='redis://127.0.0.1:6379/2'#创建默认appapp=Celery('myapp',broker=broker_url,backend=result_backend)@app.taskdefsend_sms(name,code):print("开始向%s发送验证码<%04d>"%(name,code))time.sleep(2)print("结束向%s发送验证码"%(name,))return'ok'
  1. 启动celery的主程序worker

start_worker.bat 如下

@echo off chcp 65001 >nul cd /d "%~dp0" celery -A main worker -P eventlet -l info --concurrency=4

若启动显示如下,则表示worker启动成功了

这里我们手动调用一下任务:

(.venv) >celery -A main call main.send_sms -a "[\"张三\",123]"

此时worker会有日志输出:

此时也可以通过指令查看操作结果

  1. celery其他常用命令

celery -A main status

celery -A main report# 输出:软件版本、Broker地址、结果后端、配置项等

celery -A main inspect 指令还有其他指令参数如:active, active_queues, clock, conf, memdump, memsample, objgraph, ping, query_task, registered, report, reserved, revoked, scheduled, stats

  1. celery的优雅关闭
celery -A project_name control shutdown #优雅的关闭所有的worker

control 指令还有其他的指令参数如:add_consumer, autoscale, cancel_consumer, disable_events, election, enable_events, heartbeat, pool_grow, pool_restart, pool_shrink, rate_limit, revoke, revoke_by_stamped_headers, shutdown, terminate, time_limit等

三、celery高级应用

  1. 使用独立任务目录模块,以及自动搜索,比如任务模块目录结构为:

    - main.py

    - mytasks

    - - __init__.py

    - - tasks.py #该文件名必须是这个,否则需要手动添加模块路径

    #coding: utf8'''tasks.py'''importtimefromceleryimportshared_task@shared_task(name='task_sum')deftask_sum(x,y):print("开始执行求和")time.sleep(5)print("结束执行求和")returnx+y

    此时需要修改main.py的主文件为:

    importsys,os sys.path.insert(0,os.path.dirname(os.path.abspath(__file__)))
    app.autodiscover_tasks(['mytasks'])#这里会自动搜索该目录下tasks模块下所有被shared_task修饰的任务函数

    此时重新启动worker进程,会出现一个新的任务 task_sum

    1. 自定义任务基类,用于统一处理日志,监控,重试等功能。

      #base_task.py

      # coding: utf8fromceleryimportTaskimportloggingimporttime logger=logging.getLogger(__name__)classProductionTask(Task):""" 生产环境任务基类 """max_retries=3retry_delay=60enable_retry=True# 🔧 配置:哪些异常需要重试(可被子类覆盖)retryable_exceptions=(ConnectionError,TimeoutError,OSError,# 可以添加更多)# 🔧 配置:哪些异常不重试,直接失败non_retryable_exceptions=(ValueError,TypeError,KeyError,AttributeError,# 业务逻辑错误通常不重试)def__call__(self,*args,**kwargs):"""执行任务,带监控"""task_id=self.request.idtask_name=self.name start_time=time.time()logger.info(f"[{task_name}] 开始执行, ID:{task_id}, 重试:{self.request.retries}/{self.max_retries}")try:result=super().__call__(*args,**kwargs)duration=time.time()-start_time logger.info(f"[{task_name}] 执行成功, 耗时:{duration:.2f}s")returnresultexceptExceptionase:duration=time.time()-start_time retries=self.request.retries# 🔧 判断是否应该重试should_retry=(self.enable_retryandretries<self.max_retriesandself.is_retryable_exception(e))ifshould_retry:countdown=self.retry_delay*(2**retries)logger.warning(f"[{task_name}] 执行失败:{e}, 将在{countdown}s 后重试 "f"(第{retries+1}/{self.max_retries}次)")raiseself.retry(exc=e,countdown=countdown)else:logger.error(f"[{task_name}] 执行失败:{e}, 耗时:{duration:.2f}s, "f"不满足重试条件,直接失败")raisedefis_retryable_exception(self,exc):"""判断异常是否应该重试"""# 1. 如果异常在非重试列表中,不重试ifisinstance(exc,self.non_retryable_exceptions):returnFalse# 2. 如果异常在重试列表中,重试ifisinstance(exc,self.retryable_exceptions):returnTrue# 3. 默认:不重试(保守策略)# 如果你想让默认行为是重试,可以改为 return TruereturnFalsedefon_failure(self,exc,task_id,args,kwargs,einfo):"""失败回调"""logger.error(f"任务{task_id}最终失败, 异常:{exc}, 重试次数:{self.request.retries}")

      在tasks.py文件中添加如下任务(先导入 from .base_task import ProductionTask )

      @shared_task(base=ProductionTask,bind=True,max_retries=3,retry_delay=5,name="task_send_email")deftask_send_email(self,to_email,content):"""发送邮件任务 - 自定义重试参数"""importrandomifrandom.random()<0.3:# 30%概率失败raiseConnectionError("邮件服务器暂时不可用")print(f"发送邮件到{to_email}")returnf"邮件已发送到{to_email}"
      celery -A main call task_send_email --kwargs="{\"to_email\": \"test.com\", \"content\": \"hello\"}"

  2. 带进度条的任务调度

    #coding: utf8importtimefromceleryimportTaskclassProgressTask(Task):"""带进度功能的基类"""defupdate_progress(self,current,total,extra_info=None):""" 更新任务进度 Args: current: 当前进度 total: 总数 extra_info: 额外信息(字典) """progress=int((current/total)*100)# 构建状态元数据meta={'current':current,'total':total,'progress':progress,'status':'PROGRESS'}ifextra_info:meta.update(extra_info)# 更新 Celery 状态self.update_state(state='PROGRESS',meta=meta)returnprogress
    #在tasks.py中添加任务@shared_task(bind=True,base=ProgressTask,name='task_process_data')deftask_process_data(self,total_items):""" 处理大量数据的任务 示例:处理 100 条记录 """task_id=self.request.idprint(f"[{task_id}] 开始处理{total_items}条数据")processed=0failed=0foriinrange(1,total_items+1):# 模拟处理每条数据time.sleep(0.5)# 实际业务中这里是真实处理逻辑# 模拟某些失败(10% 概率)ifi%10==0:failed+=1# 记录失败但继续处理extra_info={'last_error':f'第{i}条处理失败','failed':failed}else:processed+=1extra_info=None# 更新进度self.update_progress(current=i,total=total_items,extra_info=extra_info)# 每 10% 打印一次日志ifi%(total_items//10)==0:print(f"[{task_id}] 进度:{int(i/total_items*100)}%, 成功:{processed}, 失败:{failed}")print(f"[{task_id}] 处理完成!")return{'status':'completed','total':total_items,'processed':processed,'failed':failed}

    celery -A main call task_process_data --args=“[100]”

    收到任务后执行结果如下:

四、任务调度

  1. 异步执行任务

    #coding: utf8fromdatetimeimportdatetime,timedeltafrommainimportsend_sms## #异步调用# #send_sms.delay('李四', 1)# send_sms.apply_async(args=["张三", 1234],countdown=10)## 定时异步执行eta_time=datetime.now()+timedelta(seconds=20)result=send_sms.apply_async(args=["张三",1234],eta=eta_time)
  2. 同步执行任务

    send_sms.apply(args=["张三",1234])#这里同步调用
  3. 周期性任务调度

先在 tasks.py 中添加任务

@shared_task(base=ProductionTask,name='task_send_heartbeat')deftask_send_heartbeat(type='heartbeat'):"""发送心跳 - 每分钟"""print(f"[{datetime.now()}] 发送心跳:{type}")returnf"心跳发送成功:{type}"@shared_task(base=ProgressTask,name='task_important')deftask_important():print("我很重要!!!")return"ok"

​ 为了执行这个周期任务,我们需要在main.py中设置beat_schedule

fromcelery.schedulesimportcrontab app.conf.beat_schedule={# 任务1:每30秒执行一次'every-10-seconds':{'task':'task_send_heartbeat','schedule':timedelta(seconds=10),# 秒},# 每天 8:00 和 20:00 执行'twice-daily':{'task':'task_important','schedule':crontab(hour='8,20',minute=0),},}app.conf.timezone='Asia/Shanghai'app.conf.enable_utc=True

然后开启beat进程 start_beat.bat

@echo off chcp 65001 >nul cd /d "%~dp0" :: 激活虚拟环境 call .\.venv\Scripts\activate.bat echo Starting Celery Beat... celery -A main beat -l info pause

五、 任务监控

​ flower是celery的Web监控工具,提供了可视化界面,以及一些参数修改功能。

@echo off chcp65001>nulcd/d"%~dp0":: 激活虚拟环境 call .\.venv\Scripts\activate.batechoStarting Flower... celery-Amain flower pause

六、 在Django中使用celery详细过程

  1. 安装必要组件

    pip install celery redis==7.2.2 eventlet flower django_celery_beat
    1. 创建django项目django_celery,并新建一个app,名字为polls
INSTALL_APPS=[...'polls.apps.PollsConfig','django_celery_beat',]#一般情况下增加如下配置celeryTIME_ZONE='Asia/Shanghai'USE_TZ=True# Celery ConfigurationCELERY_BROKER_URL='redis://127.0.0.1:6379/3'CELERY_RESULT_BACKEND='redis://127.0.0.1:6379/4'CELERY_ACCEPT_CONTENT=['json']CELERY_TASK_SERIALIZER='json'CELERY_RESULT_SERIALIZER='json'CELERY_TIMEZONE=TIME_ZONE CELERY_ENABLE_UTC=USE_TZ CELERY_TASK_TRACK_STARTED=TrueCELERY_TASK_TIME_LIMIT=30*60CELERY_BEAT_SCHEDULER='django_celery_beat.schedulers:DatabaseScheduler'
  1. 在settings.py统计目录新建文件celery.py

    # myproject/celery.pyimportosfromceleryimportCeleryfromdjango.confimportsettings# 设置 Django 默认配置os.environ.setdefault('DJANGO_SETTINGS_MODULE','django_celery.settings')# 创建 Celery 应用app=Celery('myproject')# 从 Django settings 加载配置app.config_from_object('django.conf:settings',namespace='CELERY')# 自动发现任务(扫描所有 app 的 tasks.py)app.autodiscover_tasks()@app.task(bind=True,ignore_result=True)defdebug_task(self):"""调试任务"""print(f'Request:{self.request!r}')
  2. 在polls目录下新建tasks.py。

    # polls/tasks.pyimportloggingfromceleryimportshared_taskfromdjango.utilsimporttimezone logger=logging.getLogger(__name__)# ==================== Celery 任务示例 ====================@shared_taskdefsend_vote_notification(poll_id,choice_id,username):""" 投票后发送通知(Celery 异步任务) """from.modelsimportPoll,Choicetry:poll=Poll.objects.get(id=poll_id)choice=Choice.objects.get(id=choice_id)# 模拟发送通知logger.info(f" [Celery任务] 发送投票通知")logger.info(f" 用户:{username}")logger.info(f" 投票:{poll.title}")logger.info(f" 选项:{choice.text}")logger.info(f" 时间:{timezone.localtime()}")# 模拟耗时操作(展示异步效果)importtime time.sleep(3)# 模拟发送邮件耗时returnf"通知已发送给{username}"exceptExceptionase:logger.error(f"发送通知失败:{e}")raise@shared_taskdefupdate_poll_statistics(poll_id):""" 更新投票统计(Celery 异步任务) """from.modelsimportPolltry:poll=Poll.objects.get(id=poll_id)total=poll.total_votes()logger.info(f"📊 [Celery任务] 更新投票统计")logger.info(f" 投票:{poll.title}")logger.info(f" 总票数:{total}")logger.info(f" 时间:{timezone.now()}")returnf"统计已更新:{total}票"exceptExceptionase:logger.error(f"更新统计失败:{e}")raise# ==================== 额外:测试 Celery 的任务 ====================@shared_taskdeftest_task(message):""" 测试 Celery 是否正常工作 """logger.info(f"🧪 [Celery测试任务]{message}")returnf"测试成功:{message}"@shared_taskdefadd_numbers(x,y):""" 简单的加法测试 """result=x+y logger.info(f"🔢 [Celery计算任务]{x}+{y}={result}")returnresult
    1. 在manage.py统计目录开启worker,新建文件start_worker.bat
    @echo off chcp 65001 >nul cd /d "%~dp0" :: 激活虚拟环境 call .\.venv\Scripts\activate.bat set DJANGO_SETTINGS_MODULE=django_celery.settings :: Windows 下用 eventlet 池 echo Starting Celery Worker... celery -A django_celery worker -P eventlet -l info --concurrency=4 pause
    1. 可以使用我们之前学习过的任务逻辑测试命令进行任务调试
    celery -A django_celery call django_celery.celery.debug_task
    1. 在polls的vote请求成功后调用

      #polls.views.pydefvote(request,poll_id):...ifrequest.method=='POST':choice=get_object_or_404(Choice,id=choice_id,poll=poll)# ✅ 简单更新票数(允许重复投票)choice.votes+=1choice.save()send_vote_notification.delay(poll.id,choice.id,request.user.username)#实现对通知的异步转发returnredirect('polls:poll_result',poll_id=poll_id)...

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询