ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Celery 从百万到千万级:任务积压根因与调优实战

Celery 从百万到千万级:任务积压根因与调优实战 我到现在还记得那个凌晨两点十七分手机报警震醒我的瞬间。屏幕上显示的不是 CPU 飙高也不是内存挂了而是一行刺眼的字任务队列积压 800 万。我揉了揉眼睛确认不是错觉打开 Celery 的监控面板看到的是一整片任务雪崩式堆积消费者侧的 worker 一个个像被堵死的下水道——只能进不能出。那一刻我才意识到用了快两年的 Celery 分布式任务队列原来我根本还没真正理解它。这件事之后我开始从“能跑就行”转向“能扛得住才算数”一步步把任务队列从百万级日任务量干到了千万级并发吞吐。整个过程踩的坑足够写成一本反面教材。这篇就当是把我的血泪史翻出来给正在用 Celery、或者正准备上 Celery 的朋友做个参考——尤其是那些跟我一样初期觉得“并发不够就加机器”就行的人真的不是这么回事。1. 那个午夜任务积压了 800 万1.1 事故现象先还原下当时的环境。业务是一个内容平台的异步处理集群主要跑视频转码、图片压缩、数据推送、特征计算这些活。Celery 用的是 Redis 作为 brokerworker 部署在 20 台 4 核 8G 的容器里每个节点起了 8 个并发进程默认配置几乎没怎么调过。事故的直接诱因是一场周年庆活动——凌晨批量推送消息外加用户大量上传视频瞬间把队列塞满了。我原本以为 celery 的队列就是“先进先出慢慢消化”结果从积压率看消费速度根本不是缓慢而是几乎停摆。当时的表现很有意思积压数量从 0 涨到 800 万只用了不到半小时worker 的 CPU 却只有 20%。没错机器没满任务就是不走。后来排查才发现这根本不是容量不够的问题而是 Celery 默认行为把 worker 给“卡死”了。1.2 积压不等于并发不够当时团队里不止一个人提出“加机器”的方案。但加机器只能解决“并发不够”的场景可你 CPU 才 20%加了机器就是浪费。真正的瓶颈往往藏在这些地方一个 worker 进程同时只能跑一个任务但如果每个任务里都是耗时操作——比如请求外部接口、处理大文件、死等数据库锁——那么这个进程就被“占住”了后面的任务排着队也只能干等。Celery 默认会对 broker 里的任务做 prefetch预取每个 worker 进程会一次性取走一批任务到本地内存。如果这批任务里有好几个慢任务后面所有任务都只能等它们跑完。大量失败任务自动重试重试又打到同一个队列形成“重试风暴”新任务被堵在后面。这些组合在一起看起来就是典型的“任务积压”实际上是不合理的任务调度策略造成的“假性积压”。我花了很久才意识到Celery 不是一个“扔进去就会自动跑完”的黑盒它的调度策略直接决定了在压力下的表现。注意遇到任务积压第一反应先别是扩容。先看 worker 的 CPU 利用率、任务耗时分布、重试次数这些数据远比直观感受重要。2. 先别怪 Celery积压的根因是什么2.1 它只是“递纸条”不是“干活的人”要理解积压得先弄清楚 Celery 在整套系统里的定位。它本质上是一个任务分发框架你负责把任务写好扔给 broker消息中间件worker 从 broker 里取任务并执行。Celery 自己不干活干活的是 worker 进程。很多初学者的误区在于把 Celery 当成一个“并行加速器”以为只要写了 async 任务业务就会自动快起来。其实它更像一条流水线传送带broker负责传纸条每个工位worker 进程一次只能处理一张纸条处理完了才能拿下一张。我当时就是没搞清楚这条流水线的调度协议导致大量工位明明有空却在等一张大纸条做完。回头看这就像请了两个厨师但把全桌的菜都先塞给第一个厨师——第二个厨师只能干站着你还以为饭店不忙。2.2 积压的四个典型根因我把后来调优过程中遇到的根因总结成四类每次排查压力问题就先逐条对照慢任务占用某个 task 执行时间是平均值的几十倍占住 worker 不放。比如一个视频转码任务可能要跑 10 分钟但队列里其他任务只需要几百毫秒。如果用默认的进程模型这个转码任务会独占一个 worker 进程很久。prefetch 失控Celery 默认每个 worker 进程会预取 4 条任务。假设队列里前 4 条都是慢任务后面的快任务就全部排队等。等这 4 条跑完再预取 4 条依然可能有慢任务夹在里面快任务永远被堵。重试风暴一个下游接口抖动导致 10 万条任务失败每条任务默认重试 3 次每次重试还有延迟。这相当于瞬间给队列注入了 30 万条额外任务原本能消化的量直接翻倍。broker 变慢Redis 作为 broker 时如果任务体很大、序列化复杂或者查询结果条的频率过高Redis 的单线程模型会成为瓶颈worker 从 broker 拉任务都会卡住。这四条里面最隐秘的是 prefetch。它不像慢任务那么直观但影响面极大。我后来把worker_prefetch_multiplier从默认的 4 改成 1整个队列的吞吐量肉眼可见地涨了近一倍。3. 配置调优把 Celery 从“小作坊”拉到“工厂”3.1 worker 并发数与进程模型先讲最基础的并发数。Celery 的并发模型有 prefork多进程、eventlet、gevent协程以及 solo单进程只用于调试。默认是 prefork也是我推荐的生产模型因为它能用多进程隔离崩溃任务一个任务导致内存泄漏不至于拖垮整个 worker。并发数怎么定不是越高越好。我一般用这个公式起步节点 CPU 核数决定进程数上限。一个纯 CPU 密集型任务并发数设为核数即可如果是 IO 密集型大量网络请求、数据库等待可以设为核数的 2 到 4 倍。每个任务的内存占用也要算进去。假设一个任务平均 100MB节点内存 8G你开 80 个进程后果就是 OOM。我当时 4 核机器开 8 个进程看起来很合理但因为任务里有大批外部 HTTP 调用IO 等待远大于计算8 个进程还是太少。后来换成了 gevent 协程池同一个节点上的有效并发瞬间涨到上百。注意gevent 不是银弹。它坑在 monkey patch 上——如果你的任务里有些 C 扩展库比如某些数据库驱动、加密库不支持协程切换反而会把整个 worker 卡死。我后面会细讲这个坑。实操心得默认配置里worker_concurrency实际上是按进程数来算的。如果使用 gevent并发数可以设置得比较高但要反复压测我这边是先在压测环境跑上 3 天看内存和响应时间再定。3.2 prefetch 与公平调度这是我最想划重点的一节。worker_prefetch_multiplier控制每个 worker 进程一次从 broker 预取多少任务默认是 4。它的本意是减少频繁访问 broker 带来的开销但在任务耗时差异大的场景下会加剧饿死现象。想象一个护士站有 5 个护士worker 进程一堆病人任务里既有感冒发烧的快任务也有需要手术两小时的慢任务。默认模式是每个护士一次性领 4 个病人进病房哪怕后面有很多感冒病人在排队护士也得先把病房里的手术做完。而如果你把 prefetch 设为 1每个护士每次只领一个病人做完再领感冒病人很快就能被消化。调优方法很简单在启动命令里加上celery -A proj worker --concurrency16 --prefetch-multiplier1或者在配置里统一写worker_prefetch_multiplier 1注意prefetch1 不是所有场景都合适。如果任务都很小很快prefetch 设大一点能降低 broker 访问压力。我的策略是任务耗时差异大设 1任务耗时均匀且都很短设 4 或 8。实践下来差异大的场景改了之后队列积压峰值从 800 万降到几十万立竿见影。3.3 确认机制与任务丢失Celery 默认在 worker 收到任务时就向 broker 确认ack而不是执行完成后确认。这个设计的代价是如果 worker 在执行过程中崩溃任务就丢了——因为 broker 认为你已经接手并确认了。要保证“不丢任务”就要把确认时机改为“任务执行完成后”。设置方式task_acks_late True这代表 worker 完成执行后才 ack。配合worker_prefetch_multiplier1使用效果最稳任务从 broker 取出来一旦 worker 崩溃因为没 ack任务会重新投递给其他 worker。但这里有个隐藏坑acks_lateTrue时如果任务永远不结束死循环这个任务会一直占着 worker永不释放。所以必须配合超时时间task_time_limit 600 # 硬超时到了直接杀掉进程 task_soft_time_limit 540 # 软超时抛 SoftTimeLimitExceeded 异常有次我们线上有个任务因为第三方接口一直不返回默认没有超时设置导致 worker 进程成了僵尸队列越积越多。加了超时之后这个问题才算彻底根治。3.4 超时、重试与熔断任务重试是个双刃剑。Celery 的task_retry机制好用但滥用会导致重试风暴。我后来给所有可能失败的任务加了两个约束最大重试次数不超过 3 次。重试延迟使用指数退避并且加随机抖动。task(bindTrue, max_retries3, default_retry_delay60, autoretry_for(NetworkError, TimeoutError), retry_backoffTrue, retry_backoff_max600, retry_jitterTrue) def push_data(self, payload): try: api.push(payload) except ApiException as e: raise self.retry(exce, countdownrandom.randint(30, 60))指数退避的数学逻辑很简单例如第一次 30 秒第二次 60 秒第三次 120 秒。抖动是为了避免所有任务在同一时刻重试形成新的峰值。重试风暴往往比原始故障更致命因为它会把你系统的负载放大数倍。我见过最惨的一次线上数据库抖动某个队列 20 万任务失败默认重试策略是立即重试结果 5 分钟产生了几百万次重试直接把数据库再次拖垮形成了雪崩。熔断比重试更重要——当下游已经倒下了再试多少次都没用。4. 架构演进从单队列到多队列、优先级与分区4.1 多队列与路由规则任务积压事故处理完我发现另一个问题所有任务都混在同一个默认队列里视频转码这种“重任务”和发短信这种“轻任务”互相挤兑。想让轻任务优先处理光靠 prefetch 不够最直接的方式是物理拆分队列。Celery 支持多个队列每个队列对应独立命名然后用路由规则分发任务task_queues ( Queue(transcode, routing_keyvideo.transcode), Queue(notify, routing_keyuser.notify), Queue(default, routing_keytask.default), ) task_routes { tasks.video_transcode: {queue: transcode}, tasks.send_notify: {queue: notify}, }启动时分别启动不同 worker 组指向不同队列celery -A proj worker -Q transcode --concurrency8 -n transcode_worker%h celery -A proj worker -Q notify --concurrency4 -n notify_worker%h这样视频转码再慢也不会堵住通知任务。更重要的是你可以针对不同队列配置不同的并发数、超时和重试策略。比如通知任务快速失败即可转码任务可以容忍长时间执行。我还建立了一个“次品队列”专门接收那些重试了几次还是失败的任务。它们不会直接死掉而是进入另一个低优先级队列慢慢消化或者转人工处理。这比直接丢进死信队列更稳健。4.2 优先级队列的陷阱业务方经常提需求“这两类任务要优先跑”。天真地以为 Celery 设置priority参数就行其实没那么简单。在 Redis 作为 broker 时Celery 的优先级是“伪优先级”它内部用多个 list 模拟队列优先级高的任务往不同的 list 里塞消费时会优先读高优先级 list。但问题是如果高优先级队列一直有任务低优先级任务可能被饿死。RabbitMQ 对优先级的支持更可靠但数量级上也有上限最多 255 个优先级实际建议 10 个以内。我的建议是需求上真正能做到“绝对优先”的任务就别和普通任务共用队列用物理隔离的队列加独立 worker。优先级参数更适合那些“希望尽量优先但不要求严格”的场景。实操心得优先级的粒度要小于等于队列口数量。与其把 10 类任务分成 10 个优先级不如把他们拆到 3 个队列里不同队列用不同的 worker 资源和并发。优先级机制会让你陷入“调度不可控”的泥潭。4.3 任务幂等与去重千万级任务量最怕的不是重复而是“重复执行造成数据错误”。Celery 默认 At-least-once 语义配合 acks_late 时尤其如此意味着同一条任务可能被执行多次。解决重复执行的办法不是取消重试而是让任务本身幂等。业务上需要注册唯一键def process_order(order_id): # 使用 Redis setnx 加唯一执行标记 key ftask:processed:{order_id} if not redis.set(key, 1, nxTrue, ex3600): return # 已经处理过 # 真正执行任务逻辑另一个经典场景处理用户文件。如果任务重复执行可能重复发通知、重复扣费。我这里的方案是给每条任务带上业务唯一 ID在数据库里加唯一索引执行前先插入一条task_exec_log冲突就跳过。幂等设计的代价是额外判断开销但换来的是重试、故障转移、多 worker 抢占都能安心。没有幂等性千万级并发就是千万级事故。4.4 水平扩展与容量规划当单机并发顶到极限后接下来就是加机器。但加机器不是简单地“多拉几台 worker 就行”需要考虑几个问题worker 实例无状态吗任务会不会落到不同 worker 导致本地缓存失效有状态的任务尽量把状态放 Redis 或数据库。队列分区了吗单队列的 Redis 如果成了瓶颈需要把队列拆散到不同 broker 实例。我曾用一致性哈希把任务路由到不同 Redis 分片从而避免单个 Redis 实例连接数过多。自动扩容怎么做我推荐基于队列长度的指标来做。比如积压数超过阈值就触发 Kubernetes HPA 增加 worker 副本。Celery 生态里 KEDA 支持基于 Redis 或 RabbitMQ 队列长度扩容这是我现在的标准方案。当时的容量规划大致公式是日均任务量 1000 万若峰值是均值的 5 倍按 400 QPS 算单 worker 单进程单次任务平均耗时 0.2 秒则每秒能处理 5 个任务需要至少 80 个并发进程才能顶住 400QPS。保守起见再乘 1.5 到 2 倍冗余。这个估算方法比较土但很实用。5. 千万级并发下的稳定性实战5.1 监控告警你必须知道的事没有监控的任务调度系统等于闭眼开车。我早期只看着队列长度以为不增长就没事。后来发现任务“执行成功”也可能藏着延迟、重试、堆积。我现在的监控三层基础层broker 的 CPU、内存、连接数Redis 的慢日志。任务层每个队列的积压量、消费速率、平均执行时间、P99 延迟、失败率、重试率。业务层每个任务类型是否完成预期目标比如转码成功率、推送到达率。工具上Celery 自带的 Flower 能看当前 worker 状态但对历史数据不友好。后来我上了 Prometheus celery-exporter把任务指标全部打到 Grafana报警规则设了三档积压量超过 1 万warning积压量超过 10 万critical任务成功率低于 99%immediate。报警一定要带上队列名和 worker 节点否则半夜收到一个“任务失败”告警你还得一个个查日志那真是折磨。5.2 流量冲击与限流降级做千万级日任务量最怕的是业务方一次性涌入大量任务。光靠扩容有时候追不上峰值必须学会削峰填谷。我的做法是在任务入口处加一层“分诊逻辑”。任务进来先不直接进队而是按业务优先级和资源占用分类配合令牌桶限制的投递速率。Celery 本身没有内置限流但可以用两个手段在调用apply_async前用 Redis 令牌桶判断是否允许入队。使用 broker 的 QOS 限制例如 Redis 队列设定最大长度超出部分直接丢弃或落入“兜底表”等峰值过后再投递。限流不是拒绝业务而是保护整个系统的可用性。曾有一次渠道方忽然回传了大量待处理数据我们全靠入口限流兜住否则下游数据库又要被压垮。5.3 慢任务治理与拆分千万级规模下只要存在 0.01% 的慢任务也会给你制造海量的尾部延迟。我这里的慢任务治理分三步第一步找出拖后腿的任务。通过任务执行时间分布筛出平均耗时超阈值或者 P99 明显上升的任务类型。第二步拆分任务。一个任务里做了太多事就把它拆成多个小任务。比如原来“下载视频 - 转码 - 截图 - 上传 CDN”拆成四个独立任务放到不同队列分别扩容。这样单任务执行时间大幅下降prefetch 策略也更可控。第三步大任务走独立资源池。有些任务确实没法拆比如超高清视频全量转码。就把它们放到单独的 transcode 队列用少量 worker 超长超时专门跑避免污染主处理链路。我还常用“先切片、后聚合”的模式一个大文件切成 100 片每片一个小任务处理最后汇总。这也让水平扩容的收益更明显——100 片任务可以分散到 100 台机器上并行处理。5.4 broker 选型与性能对比很多人问我Redis 够不够用说实话在千万级日任务量下Redis 仍然可以用但你必须接受它的两个短板它是单一内存存储队列达到几百万条时内存占用很大。它是单线程模型出队入队大量并发时CPU 会成为瓶颈。如果你用 Redis 作为 broker我建议任务体尽量轻量只放 ID 和必要参数详细信息通过任务执行时从 DB/缓存取。开启visibility_timeoutSQS 的概念在 Redis 上不算原生但可以通过 expire 控制防止任务卡在 unacked 状态。给 broker 单独部署不要让业务缓存和队列共用 Redis 实例。而 RabbitMQ 更适合大量持久化任务场景它天然支持 ACK、优先级、死信队列性能在持久化模式下也很稳。但运维成本比 Redis 高。我现在的方案是默认走 Redis重要且不允许丢失的任务走 RabbitMQ两条链路隔离。下面是我压测过的大致对比供参考维度RedisRabbitMQ入队吞吐很高但受单线程限制高支持多线程写入持久化弱依赖 RDB/AOF强消息可持久化重试/ACK较弱需要靠 Celery 层原生支持行为可靠运维成本低中高优先级伪优先级原生优先级适用场景海量短任务允许偶发丢失交易类、订单类、不可丢任务6. 血泪踩坑实录常见问题排查速查6.1 任务“神秘消失”现象任务在队列里存在worker 也显示接收了但执行日志里没有输出最后任务不见了。最常见的原因就是acks_lateFalse时 worker 进程崩溃任务已经 ack所以不会重发。检查时先看 worker 日志里有没有Hard time limit或WorkerLostError。排查步骤查看 task 的启动时间和结束时间。看 worker 进程是否有 OOM killed 记录。确认task_acks_on_failure_or_timeout的配置默认 True失败超时也会 ack导致不重试。如果希望超时或失败后不丢任务设置task_acks_on_failure_or_timeout False但注意这会增加重试量。需要结合“任务幂等”一起使用否则重复处理的副作用会让你更头疼。6.2 worker 卡死与内存泄漏我试过 gevent 池一度很爽直到一个任务里调用了subprocess等待外部进程导致整个 worker 全部挂起。排查时发现协程被卡在没有被 monkey patch 的 C 扩展上。另一个常见问题是任务里用了大量本地变量或引用了框架全局对象导致内存只升不降。我的做法是对 worker 进程设置ulimit防止单进程内存无限增长。部署后监控各 worker 进程的 RSS。发现内存持续上涨先通过celery -A proj purge清掉积压任务再逐步找出罪魁祸首任务。某些任务即使抛出异常也会占住内存不释放。这时只能在配置里设置maxtasksperchild让每个 worker 子进程在处理预处理任务后主动重启worker_max_tasks_per_child 200 worker_max_memory_per_child 300000 # 300M这是非常土但非常有效的保命手段。6.3 Redis 连接被占满Celery worker 每个进程默认会建一条 broker 连接如果你开 100 个进程那 Redis 侧至少有 100 条连接。如果并发再叠加连接数很容易打满 maxclients。解决方案开启连接池复用设置broker_pool_limit。使用redis作为 broker 时为 Celery 单独部署 Redis不要把其他业务连接混在一起。拉大worker_prefetch_multiplier减少 worker 频繁轮询 broker 的次数。我踩过的具体表现是Redis 连接数到达阈值新的 worker 起不来旧 worker 偶尔报ConnectionError。检查命令redis-cli info clients如果 connected_clients 接近 maxclients就该拆分 Redis 或调低 worker 进程数。6.4 重复执行的噩梦用了acks_lateTrue和worker_prefetch_multiplier1后任务重复的概率明显上升——因为这本来就是为了“不丢任务”付出的代价。尤其是网络抖动时worker 与 broker 断开broker 不知道该任务已经被取走还是没取走会重新投递。要避免重复造成业务错误我已经在 4.3 里写过了——幂等设计是唯一根治路径。判断一条任务是否执行过不要在内存里判断要用 Redis 或数据库的原子操作。比如# 伪代码用任务唯一ID做去重 key fcelery:dedup:{task_id} if not redis.set(key, 1, ex3600, nxTrue): return注意去重时间窗口要比任务最长执行时间大否则窗口过期后任务再次执行还是会重复。7. 最后再分享几个实操经验如果你也在用 Celery我个人强烈建议从一开始就考虑下面这些事而不是等线上出事再补救永远不要让任务直接操作“不好惹”的第三方接口。第三方超时不可控一定要包一层“超时 重试 熔断”。线上改配置别一次改太多参数。我每次只改一项然后观察 24 小时。有次同时改了acks_late、prefetch、并发模型出问题时根本不知道是哪一项导致的。给任务命名要规范。到千万级后你会天天看监控面板。任务名一乱排查效率下降一半。定期做故障演练。故意停掉几个 worker看队列积压会不会自动恢复故意让 Redis 重启看有没有任务丢失。没有演练过的系统永远不知道它在故障下有多脆弱。还有一个我特别想说的Celery 的依赖环境很容易“脏”。有些任务库需要特定版本的 Python有些 worker 节点缺底层依赖导致任务一执行就崩溃。我给每个任务队列容器化了 worker 镜像这样扩容出来的机器行为完全一致省掉了大量“我这里没问题”的扯皮。走到今天我仍然不觉得 Celery 是完美的它有很多“约定高于配置”的默认行为对新手并不友好。但只要你理解它背后的调度逻辑按队列拆、按任务分、按幂等兜底、按监控观测千万级并发并不是遥不可及的事。至少我这样的普通人靠踩坑和持续迭代也就熬过来了。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进