Windows Celery进阶之路:自定义基类、进度监控、周期任务与Django全整合

发布时间:2026/8/11 14:16:54
Windows Celery进阶之路:自定义基类、进度监控、周期任务与Django全整合 一、celery 介绍Celery是一个简单的、快速的灵活且可靠的分布式系统用于处理大量消息同时提供了一些工具来维护这样的一个系统。这是一个专注于实时处理的任务队列同时也支持任务调度。Celery 支持自主配置消息队列结果存储并发序列化等这里我们使用在windows下使用redis作为消息队列和结果存储使用eventlet并发, 序列化使用json的方式。二、celery基本应用安装 Python 及相关组件pip install celery redis7.1.3eventlet flower新建一个主应用文件名我们命名为 main.pyimporttimefromceleryimportCelery broker_urlredis://127.0.0.1:6379/1result_backendredis://127.0.0.1:6379/2#创建默认appappCelery(myapp,brokerbroker_url,backendresult_backend)app.taskdefsend_sms(name,code):print(开始向%s发送验证码%04d%(name,code))time.sleep(2)print(结束向%s发送验证码%(name,))returnok启动celery的主程序workerstart_worker.bat 如下echo off chcp 65001 nul cd /d %~dp0 celery -A main worker -P eventlet -l info --concurrency4若启动显示如下则表示worker启动成功了这里我们手动调用一下任务(.venv) celery -A main call main.send_sms -a [\张三\,123]此时worker会有日志输出此时也可以通过指令查看操作结果celery其他常用命令celery -A main statuscelery -A main report# 输出软件版本、Broker地址、结果后端、配置项等celery -A main inspect 指令还有其他指令参数如active, active_queues, clock, conf, memdump, memsample, objgraph, ping, query_task, registered, report, reserved, revoked, scheduled, statscelery的优雅关闭celery -A project_name control shutdown #优雅的关闭所有的workercontrol 指令还有其他的指令参数如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高级应用使用独立任务目录模块以及自动搜索比如任务模块目录结构为- main.py- mytasks- - __init__.py- - tasks.py #该文件名必须是这个否则需要手动添加模块路径#coding: utf8tasks.pyimporttimefromceleryimportshared_taskshared_task(nametask_sum)deftask_sum(x,y):print(开始执行求和)time.sleep(5)print(结束执行求和)returnxy此时需要修改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自定义任务基类用于统一处理日志监控重试等功能。#base_task.py# coding: utf8fromceleryimportTaskimportloggingimporttime loggerlogging.getLogger(__name__)classProductionTask(Task): 生产环境任务基类 max_retries3retry_delay60enable_retryTrue# 配置哪些异常需要重试可被子类覆盖retryable_exceptions(ConnectionError,TimeoutError,OSError,# 可以添加更多)# 配置哪些异常不重试直接失败non_retryable_exceptions(ValueError,TypeError,KeyError,AttributeError,# 业务逻辑错误通常不重试)def__call__(self,*args,**kwargs):执行任务带监控task_idself.request.idtask_nameself.name start_timetime.time()logger.info(f[{task_name}] 开始执行, ID:{task_id}, 重试:{self.request.retries}/{self.max_retries})try:resultsuper().__call__(*args,**kwargs)durationtime.time()-start_time logger.info(f[{task_name}] 执行成功, 耗时:{duration:.2f}s)returnresultexceptExceptionase:durationtime.time()-start_time retriesself.request.retries# 判断是否应该重试should_retry(self.enable_retryandretriesself.max_retriesandself.is_retryable_exception(e))ifshould_retry:countdownself.retry_delay*(2**retries)logger.warning(f[{task_name}] 执行失败:{e}, 将在{countdown}s 后重试 f(第{retries1}/{self.max_retries}次))raiseself.retry(exce,countdowncountdown)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(baseProductionTask,bindTrue,max_retries3,retry_delay5,nametask_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\}带进度条的任务调度#coding: utf8importtimefromceleryimportTaskclassProgressTask(Task):带进度功能的基类defupdate_progress(self,current,total,extra_infoNone): 更新任务进度 Args: current: 当前进度 total: 总数 extra_info: 额外信息字典 progressint((current/total)*100)# 构建状态元数据meta{current:current,total:total,progress:progress,status:PROGRESS}ifextra_info:meta.update(extra_info)# 更新 Celery 状态self.update_state(statePROGRESS,metameta)returnprogress#在tasks.py中添加任务shared_task(bindTrue,baseProgressTask,nametask_process_data)deftask_process_data(self,total_items): 处理大量数据的任务 示例处理 100 条记录 task_idself.request.idprint(f[{task_id}] 开始处理{total_items}条数据)processed0failed0foriinrange(1,total_items1):# 模拟处理每条数据time.sleep(0.5)# 实际业务中这里是真实处理逻辑# 模拟某些失败10% 概率ifi%100:failed1# 记录失败但继续处理extra_info{last_error:f第{i}条处理失败,failed:failed}else:processed1extra_infoNone# 更新进度self.update_progress(currenti,totaltotal_items,extra_infoextra_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]”收到任务后执行结果如下四、任务调度异步执行任务#coding: utf8fromdatetimeimportdatetime,timedeltafrommainimportsend_sms## #异步调用# #send_sms.delay(李四, 1)# send_sms.apply_async(args[张三, 1234],countdown10)## 定时异步执行eta_timedatetime.now()timedelta(seconds20)resultsend_sms.apply_async(args[张三,1234],etaeta_time)同步执行任务send_sms.apply(args[张三,1234])#这里同步调用周期性任务调度先在 tasks.py 中添加任务shared_task(baseProductionTask,nametask_send_heartbeat)deftask_send_heartbeat(typeheartbeat):发送心跳 - 每分钟print(f[{datetime.now()}] 发送心跳:{type})returnf心跳发送成功:{type}shared_task(baseProgressTask,nametask_important)deftask_important():print(我很重要)returnok​ 为了执行这个周期任务我们需要在main.py中设置beat_schedulefromcelery.schedulesimportcrontab app.conf.beat_schedule{# 任务1每30秒执行一次every-10-seconds:{task:task_send_heartbeat,schedule:timedelta(seconds10),# 秒},# 每天 8:00 和 20:00 执行twice-daily:{task:task_important,schedule:crontab(hour8,20,minute0),},}app.conf.timezoneAsia/Shanghaiapp.conf.enable_utcTrue然后开启beat进程 start_beat.batecho 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 chcp65001nulcd/d%~dp0:: 激活虚拟环境 call .\.venv\Scripts\activate.batechoStarting Flower... celery-Amain flower pause六、 在Django中使用celery详细过程安装必要组件pip install celery redis7.2.2 eventlet flower django_celery_beat创建django项目django_celery并新建一个app名字为pollsINSTALL_APPS[...polls.apps.PollsConfig,django_celery_beat,]#一般情况下增加如下配置celeryTIME_ZONEAsia/ShanghaiUSE_TZTrue# Celery ConfigurationCELERY_BROKER_URLredis://127.0.0.1:6379/3CELERY_RESULT_BACKENDredis://127.0.0.1:6379/4CELERY_ACCEPT_CONTENT[json]CELERY_TASK_SERIALIZERjsonCELERY_RESULT_SERIALIZERjsonCELERY_TIMEZONETIME_ZONE CELERY_ENABLE_UTCUSE_TZ CELERY_TASK_TRACK_STARTEDTrueCELERY_TASK_TIME_LIMIT30*60CELERY_BEAT_SCHEDULERdjango_celery_beat.schedulers:DatabaseScheduler在settings.py统计目录新建文件celery.py# myproject/celery.pyimportosfromceleryimportCeleryfromdjango.confimportsettings# 设置 Django 默认配置os.environ.setdefault(DJANGO_SETTINGS_MODULE,django_celery.settings)# 创建 Celery 应用appCelery(myproject)# 从 Django settings 加载配置app.config_from_object(django.conf:settings,namespaceCELERY)# 自动发现任务扫描所有 app 的 tasks.pyapp.autodiscover_tasks()app.task(bindTrue,ignore_resultTrue)defdebug_task(self):调试任务print(fRequest:{self.request!r})在polls目录下新建tasks.py。# polls/tasks.pyimportloggingfromceleryimportshared_taskfromdjango.utilsimporttimezone loggerlogging.getLogger(__name__)# Celery 任务示例 shared_taskdefsend_vote_notification(poll_id,choice_id,username): 投票后发送通知Celery 异步任务 from.modelsimportPoll,Choicetry:pollPoll.objects.get(idpoll_id)choiceChoice.objects.get(idchoice_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})raiseshared_taskdefupdate_poll_statistics(poll_id): 更新投票统计Celery 异步任务 from.modelsimportPolltry:pollPoll.objects.get(idpoll_id)totalpoll.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): 简单的加法测试 resultxy logger.info(f [Celery计算任务]{x}{y}{result})returnresult在manage.py统计目录开启worker新建文件start_worker.batecho off chcp 65001 nul cd /d %~dp0 :: 激活虚拟环境 call .\.venv\Scripts\activate.bat set DJANGO_SETTINGS_MODULEdjango_celery.settings :: Windows 下用 eventlet 池 echo Starting Celery Worker... celery -A django_celery worker -P eventlet -l info --concurrency4 pause可以使用我们之前学习过的任务逻辑测试命令进行任务调试celery -A django_celery call django_celery.celery.debug_task在polls的vote请求成功后调用#polls.views.pydefvote(request,poll_id):...ifrequest.methodPOST:choiceget_object_or_404(Choice,idchoice_id,pollpoll)# ✅ 简单更新票数允许重复投票choice.votes1choice.save()send_vote_notification.delay(poll.id,choice.id,request.user.username)#实现对通知的异步转发returnredirect(polls:poll_result,poll_idpoll_id)...

相关新闻