发布时间:2026/8/11 14:16:55
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)...

相关新闻

2026/8/11 14:16:55

Android Jetpack核心组件详解与实战技巧

1. Android Jetpack核心组件全景解析 作为Google官方推出的Android开发组件集合,Jetpack已经成为现代Android应用开发的基石。我在过去三年主导过7个大型商业项目的架构设计,深刻体会到合理运用Jetpack组件能提升40%以上的开发效率。下面将结合典型应用场…

2026/8/11 14:11:55

7个关键技巧让你快速掌握免费2D CAD软件LibreCAD

7个关键技巧让你快速掌握免费2D CAD软件LibreCAD 【免费下载链接】LibreCAD LibreCAD is a cross-platform 2D CAD program. It can read DXF/DWG, and write DXF/DWG/PDF/SVG files. It supports point/line/circle/ellipse/parabola/hyperbola/spline primitives. The GUI is…

2026/8/11 16:12:00

Linux文件系统访问原理与性能优化实战

1. 项目概述:Linux文件系统访问的核心逻辑 在Linux系统中,文件系统如同一个精密的图书馆管理系统。它不仅要管理文件的存储位置,还要处理权限控制、数据缓存、磁盘调度等复杂事务。与Windows的盘符划分不同,Linux采用单一的树状结…

2026/8/11 16:12:00

Hermes 一体化 Windows 资源包,内置运行组件省去环境调试

Windows 本地运行 Hermes 智能 Agent,整合资源包简化搭建流程 很多想要体验 Hermes Agent 的用户都会碰到难题,原生部署需要处理大量环境配置工作,各类依赖安装、路径调整步骤繁琐,部署过程中容易接连出现报错,影响正…

2026/8/11 16:12:00

AI编程助手剪贴板集成:解决终端多模态输入难题的技术方案

1. 项目概述:当AI助手遇上“粘贴”的尴尬 如果你和我一样,日常重度依赖Claude Code或DeepSeek这类AI编程助手,那你一定遇到过这个让人瞬间血压升高的场景:在终端里,你习惯性地按下 CtrlV ,想把刚刚截取的…

2026/8/11 16:12:00

Telegram MTProto 协议,对比主流 IM 协议强在哪

对比对象:Signal 双棘轮协议、Matrix‑Olm/Megolm、传统 XMPP。MTProto 是一套完整传输 加密的自研协议,不只是单纯端到端加密组件。✅ MTProto 相对突出优势移动弱网性能优秀 不依赖标准 TLS,自定义二进制帧、轻量化握手、智能重连&#xf…

2026/8/11 16:12:00

单片机毕设项目:基于单片机的 LCD 实时显示燃气安防监测终端设计 基于 51 单片机继电器驱动排风水泵安防控制系统(017602)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于嵌入式单片机,Java、小程序技术领域和毕业项目实战 ✌️…

2026/8/11 16:07:00

绿色创新绩效综合测算数据

时间跨度2011-2024 年区域跨度中国大陆全部 31 个省级行政区数据格式Excel形式数据简介本数据为中国省级区域绿色创新绩效综合指标体系测算数据,基于中国各省份统计年鉴、《中国统计年鉴》、《中国环境统计年鉴》等多源资料整理而成。区域绿色创新绩效是指在绿色创新…

2026/8/11 3:03:40

如何快速生成中国车牌图片:Python开源工具完整指南

如何快速生成中国车牌图片:Python开源工具完整指南 【免费下载链接】chinese_license_plate_generator 中国车牌生成器 项目地址: https://gitcode.com/gh_mirrors/ch/chinese_license_plate_generator 中国车牌生成器是一个基于Python的开源项目&#xff0c…

2026/8/11 5:34:14

当 LLM 遇见大文档:主流开源项目如何处理上下文超限

从 Agentic Loop 到 Repo Map,七种策略与六类陷阱引言:128K vs 10MB 的硬冲突 2026 年的 LLM 上下文窗口已达到 128K ~ 1M token(≈ 0.5MB ~ 4MB 文本),但 LLM 想要处理的真实数据规模远远超过这个量级:真实…

2026/8/11 0:00:39

前后端分离项目中控制台与接口工具数据差异排查指南

1. 问题现象解析:控制台与Apifox的数据差异 最近在调试一个前后端分离项目时,遇到了一个典型问题:后端服务在本地开发环境控制台能正常输出查询数据,但通过Apifox测试时却返回空结果。这种"控制台有数据,接口工具…

2026/8/11 0:00:39

AI编程实战:从Claude Code踩坑到游戏开发入门

1. 从“AI能帮我做游戏”到“AI让我重新学编程”最近身边不少朋友,尤其是一些非技术背景、但对游戏开发有浓厚兴趣的朋友,都在问我同一个问题:“听说现在用Claude Code这种AI编程工具,小白也能做游戏了,是真的吗&#…

2026/8/10 11:20:30

实测才敢推 AI论文网站 2026最新测评与推荐

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。一、综…

2026/8/10 11:20:30

2026必备!AI论文网站测评:最新推荐与深度对比

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

2026/8/11 3:05:11

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…