☰
PHP 分布式任务调度实战:基于 Redis 打造轻量可靠调度器
2026/10/3 3:08:25 网站建设 项目流程

1. 为什么需要在 PHP 里实现分布式任务调度

1.1 单机定时任务为什么撑不住了

我刚接手团队这套系统时,业务的定时任务还全部挂在 Linux 的 crontab 上。每天早上两点跑报表、每天凌晨同步会员数据、每十分钟清理一次临时文件,看起来一切正常。可等业务量上来之后,问题开始扎堆出现:同一个任务部署在三四台服务器上,结果每台机器都在跑,数据同步重复执行,报表数据翻了三倍;后来把任务指定到一台机器执行,这台机器夜里一挂,第二天早上全公司都在问报表去哪了;再后来任务多了,crontab 里的记录长得像天书,谁也不敢动,动一下线上就出事。

这时候我意识到,crontab 本身并没有错,但用它做分布式环境下的任务调度,就像拿纸和笔去记库存——单机还好,一旦需要多机协作、统一管理、失败重试、动态扩展,它就完全不合适了。我们真正需要的是一个自己能掌控的调度系统:任务集中存储,多台机器竞争消费,同一个时间点同一个任务只被一台机器执行,任务失败能自动重试,执行过程可见可查。在 PHP 体系里,这些需求完全可以靠以 Redis 为核心的轻量调度器来实现。

1.2 分布式任务调度的四个核心问题

我在设计这套方案钱,明确了自己要解决的问题,其实不外乎四件事。

第一是任务分发。任务产生后,要能方便地进入一个统一的任务池,然后由多个 worker 节点各自拉取,而不是靠人工手动分配到某台服务器。这一步需要的是“任务队列”,队列里放着待执行任务的描述信息,比如任务类名、参数、优先级、执行时间。

第二是任务互斥。分布式环境下最怕的就是同一任务被多个节点同时消费。比如清理过期订单,两个 worker 同时跑,很可能把同一批订单处理两遍。要解决这个问题,需要引入分布式锁,保证同一条任务在同一时刻只有一个消费者在处理。

第三是任务时序。很多任务并不是立刻执行的,而是“延迟 5 分钟通知用户”“每天凌晨跑全量”“10 秒后再试一次”。这就要求队列里的任务不能只按 FIFO 弹出,还要支持定时、延迟、优先级等时间约束。

第四是任务可靠性。worker 执行过程中可能崩溃、任务可能抛异常、Redis 可能短时间抖动。调度器需要把这些情况考虑进去,至少保证任务不会因为一次执行失败就永久丢失,或者因为进程重启就出现大面积重复消费。

想清楚了这四点,后面的技术选型就有了方向。每种方案我都会拿实际场景去掂量,而不是看着流行就用。

2. 技术选型:为什么我选 Redis 作为调度核心

2.1 先用数据库、再用 RabbitMQ、最后落在 Redis

最开始图省事,我用 MySQL 表存任务,一个字段放任务类型,一个字段放 JSON 参数,再加一个状态字段。worker 循环去SELECT ... FOR UPDATE拉取任务,更新状态为 processing。小规模下确实能跑,但并发一上来,行锁争用严重,数据库连接池很快被打满,而且轮询查询的效率很低。后来我试过 RabbitMQ,可靠性和功能都没得说,但团队里没人懂运维,部署一套 RabbitMQ 集群对我们来说太重了。最后我把目光放到了 Redis 上——团队本来就在用,不用额外引入运维成本,而且 Redis 的 List、ZSet、Hash 这三种数据结构,简直天生就是为任务调度准备的。

List 是天然的 FIFO 队列,LPUSH投任务、BRPOP阻塞消费,生产者消费者模型一把梭。ZSet 适合做延迟队列和定时任务,member 存任务 ID,score 存执行时间戳,轮询时按时间范围取出该执行的任务。Hash 则用来存任务的元信息和执行状态,比如 key 是任务 ID,field 是执行次数、下次执行时间、最后错误信息等。这三种结构组合在一起,几乎不用额外的中间件,就把调度器的核心存储解决了。

补充一句:如果你的系统后续需要非常复杂的分流、死信、消息回溯,RabbitMQ 或 Kafka 依然是更专业的方案。本文这套是基于 Redis 的轻量级做法,适合大多数中小规模业务,同时也特别适合只想用 PHP 解决调度问题的团队。

2.2 Redis 数据结构在调度场景里的具体角色

我用一个具体例子说明:假设要做一个“过期未支付订单 30 分钟后自动关闭”的功能。

订单创建后,把订单关闭任务写进 Redis。任务内容封装成一个 JSON,里面带着任务类型order_close、订单 ID、创建时间。这个任务不会立刻被消费,而是要等 30 分钟之后才执行。我把它塞进 ZSet,member 是任务唯一 ID,score 是time() + 1800。然后有一个调度进程,每秒去 ZSet 里ZRANGEBYSCORE拉取 score 小于等于当前时间的前 100 个任务,再把它们推到 List 队列里,让 worker 真正执行。

这套结构里,ZSet 就是“等待时间到达的缓冲区”,List 是“准备执行的待消费队列”。调度进程负责把时间到期的任务从 ZSet 挪到 List,worker 只管从 List 弹任务执行,职责非常清晰。后面我会给出完整代码。

2.3 分布式环境下如何用 Redis 锁保证任务互斥

既然多台机器上是多个 worker 进程并发消费同一个 List,那我就要保证某个任务在处理期间,其他机器上的 worker 不会碰它。我先用的是 Redis 的原子操作SETNX加锁。每一个任务被 worker 从 List 弹出时,worker 先尝试给任务 ID 设置一个锁,比如SET lock:taskid taskid NX EX 300,拿到锁才执行,拿不到锁就说明已经有别的 worker 在处理,直接舍弃这条任务。

锁的过期时间也很讲究。设置太短,任务执行超过锁时间,另一个 worker 就会再次拿到锁,导致重复执行;设置太长,万一执行进程崩溃,这个任务会被锁到天荒地老。我的经验是给锁设置一个比预估执行时间长 3 到 5 倍的过期时间,然后任务执行完成后主动删除锁。还有一个容易忽略的点是,删除锁时不要只DEL key,因为可能锁已经过期被其他进程重新拿到了,这时候你删除的就是别人的锁,会造成连锁混乱。正确做法是用 Lua 脚本,先判断 value 是否是自己的标识,再执行删除。这段代码我放在后面实操部分。

2.4 为什么不直接用 Swoole 或 Workerman

我知道很多人一听 PHP 做常驻任务,第一反应就是 Swoole 或者 Workerman。确实,这两个框架能极大提升 PHP 对长连接、并发 IO 的处理能力,Workerman 本身也带Timer和AsyncTcpConnection,可以实现进程内定时任务。但我最终没有直接把它们作为调度器主体,原因是:

  • Workerman 的Timer是单进程内的定时器,如果调度器部署在多台机器上,还是需要借助外部存储来做任务的统一管理;
  • Swoole 的Process、Coroutine特性很强,但团队学习成本不小;
  • 我们需要的调度器本质上是一个“基于时间轴的任务仓库 + 动态扩容的消费进程”,用纯 PHP 配合 Redis 和 POSIX 进程控制,就能做得足够稳,还能完全掌控每一行代码。

于是我的最终方案是:调度逻辑用纯 PHP CLI 脚本实现,常驻内存用 Supervisor 托管,进程间通信用 Redis,执行任务用 PHP 的子进程或异步方式。这样的好处是任何能跑 PHP 的环境都能复现,不依赖框架,也方便迁移。

3. 从 0 搭一个分布式任务调度器

我会把这套实现拆成四个部分:任务实体定义、任务投递、任务存储与调度、任务消费。所有代码都经过线上验证,核心逻辑可以直接抄,但建议读者根据自己框架的命令行组件做二次适配。

3.1 整体架构与模块拆分

先画一张逻辑视图(虽然我画不了图,但用文字描述):Producer(业务代码)调用TaskProducer::push(),把任务写入 Redis 的 ZSet 中。Scheduler(调度进程)每秒钟执行一次,从 ZSet 中拉取到期任务,转写到 List。Consumer(多个 worker 进程)从 List 中用BRPOP阻塞弹出任务,拿分布式锁后执行任务,根据执行结果更新任务状态、安排重试或记录死信。Consumer 可以部署在同一台机器的多个进程,也可以分布在多台服务器上,只要它们连的是同一个 Redis。

模块名称我起得比较朴素,你根据自己的编码习惯改就行。

模块职责技术要点
TaskProducer业务侧投递任务封装任务参数、编号、延迟时间
TaskEntity任务实体任务类型、参数、重试次数、状态
TaskScheduler定时调度ZSet 时间轮询、批量迁移到 List
TaskWorker任务消费分布式锁、执行任务、失败重试
TaskWatcher监控辅助任务长度、失败数、耗时上报

3.2 代码实现:任务投递模块

我习惯先写一个任务实体类,方便所有地方统一引用字段名。实际项目里可能会关联你当前框架的 Model,但用纯数组也完全够用。下面这段是简化后的代码:

<?php declare(strict_types=1); class TaskEntity { public const STATUS_PENDING = 0; public const STATUS_EXECUTING = 1; public const STATUS_SUCCESS = 2; public const STATUS_FAILED = 3; public function __construct( public string $id, public string $type, public array $params, public int $delay = 0, public int $maxRetry = 3, public int $status = self::STATUS_PENDING, public int $retryCount = 0, public int $nextRunAt = 0, public ?string $lastError = null, ) {} public function toArray(): array { return get_object_vars($this); } public static function fromArray(array $data): self { return new self( id: $data['id'], type: $data['type'], params: $data['params'], delay: $data['delay'] ?? 0, maxRetry: $data['maxRetry'] ?? 3, status: $data['status'] ?? self::STATUS_PENDING, retryCount: $data['retryCount'] ?? 0, nextRunAt: $data['nextRunAt'] ?? 0, lastError: $data['lastError'] ?? null, ); } }

然后是一个任务投递的封装类。这里最核心的是生成一个唯一 ID,我用的方法是md5(uniqid((string) mt_rand(), true)),虽然没有 UUID 那么高级,但冲突概率在业务场景里可以忽略。如果你追求更严谨,可用ramsey/uuid库。

<?php declare(strict_types=1); class TaskProducer { private const PENDING_ZSET = 'task:pending:zset'; public function __construct(private Redis $redis) {} public function push(string $type, array $params, int $delay = 0, int $maxRetry = 3): string { $task = new TaskEntity( id: md5(uniqid((string) mt_rand(), true)), type: $type, params: $params, delay: $delay, maxRetry: $maxRetry, nextRunAt: time() + $delay, ); $this->redis->zAdd(self::PENDING_ZSET, $task->nextRunAt, $task->id); $this->redis->hSet('task:data:' . $task->id, 'payload', json_encode($task->toArray(), JSON_UNESCAPED_UNICODE)); return $task->id; } }

这里要注意:我把任务的完整 payload 存到了 Hash 里,而不是直接在 ZSet 里保存 JSON。原因是 ZSet 的 member 必须唯一,而且用 Hash 方便后期扩展状态字段。当然你也可以省掉 Hash 直接zAddJSON,但如果任务执行失败需要更新重试次数,频繁修改 member 就不是那么方便了。所以我在设计上让 member 只是任务 ID,详情全部走 Hash。

3.3 代码实现:调度进程 Scheduler

调度进程是一个 CLI 脚本,逻辑很简单,循环执行transfer(),每次usleep(100000)也就是 0.1 秒后继续。它可以跑在任意一台或多台服务器上,因为zRangeByScore取任务的动作是原子的,不会重复取到同一条任务,所以我们完全可以让多台机器同时跑多个调度进程,只要保证数量别太多就行。

<?php declare(strict_types=1); class TaskScheduler { private const PENDING_ZSET = 'task:pending:zset'; private const READY_LIST = 'task:ready:list'; private const LOCK_KEY = 'task:scheduler:lock'; private const LOCK_TTL = 5; public function __construct(private Redis $redis) {} public function run(): void { // 调度进程本身也加一个锁,防止多个调度进程同时搬移任务 if (!$this->redis->set(self::LOCK_KEY, '1', ['NX', 'EX' => self::LOCK_TTL])) { return; } $now = time(); $batchSize = 100; $taskIds = $this->redis->zRangeByScore(self::PENDING_ZSET, 0, $now, ['limit' => [0, $batchSize]]); if (!$taskIds) { return; } // 把到期任务从 ZSet 中移除,并推入 List foreach ($taskIds as $taskId) { $removed = $this->redis->zRem(self::PENDING_ZSET, $taskId); if ($removed) { $this->redis->lPush(self::READY_LIST, $taskId); } } } }

这里有一个隐藏的细节:我通过zRangeByScore取出一批任务后,删除和写入 List 并不是同一个原子操作。如果调度进程在遍历时崩溃,可能有任务从 ZSet 移除了但没有进 List,从而丢失。解决思路是用 Redis 事务或者 Lua 脚本把整个过程原子化。线上我建议写 Lua 脚本,这里我给出一个版本:

local taskIds = redis.call('zRangeByScore', KEYS[1], 0, ARGV[1], 'limit', 0, ARGV[2]) if #taskIds == 0 then return {} end for i, taskId in ipairs(taskIds) do redis.call('zRem', KEYS[1], taskId) redis.call('lPush', KEYS[2], taskId) end return taskIds

在 PHP 里调用:

$script = <<<'LUA' local taskIds = redis.call('zRangeByScore', KEYS[1], 0, ARGV[1], 'limit', 0, ARGV[2]) if #taskIds == 0 then return {} end for i, taskId in ipairs(taskIds) do redis.call('zRem', KEYS[1], taskId) redis.call('lPush', KEYS[2], taskId) end return taskIds LUA; $taskIds = $this->redis->eval($script, ['task:pending:zset', 'task:ready:list'], [time(), 100]);

这段脚本保证“ZSet 删除 + List 写入”要么都成功,要么都不执行。调度进程频繁运行,每次执行本身非常轻快,Redis 单实例每秒可以处理数万次这样的小脚本,不会成为瓶颈。

3.4 代码实现:Consumer 消费进程

Consumer 是真正跑任务的进程。它用BRPOP从 List 里阻塞弹出任务 ID,然后调handleTask()执行。BRPOP在没任务时会阻塞最多 N 秒,这样比普通 for 循环空转地轮询节省 CPU 得多。

<?php declare(strict_types=1); class TaskWorker { private const READY_LIST = 'task:ready:list'; private const LOCK_TTL = 300; public function __construct(private Redis $redis) {} public function run(): void { while (true) { $task = $this->redis->brPop(self::READY_LIST, 5); if (!$task) { continue; } $taskId = $task[1]; $this->process($taskId); } } protected function process(string $taskId): void { $payload = $this->redis->hGet('task:data:' . $taskId, 'payload'); if (!$payload) { return; } $task = TaskEntity::fromArray(json_decode($payload, true)); // 获取分布式锁,防止其他 worker 并发执行同一任务 $lockValue = (string) getmypid(); $locked = $this->redis->set('task:lock:' . $taskId, $lockValue, ['NX', 'EX' => self::LOCK_TTL]); if (!$locked) { return; } $task->status = TaskEntity::STATUS_EXECUTING; $this->saveTask($task); try { $result = $this->executeTask($task); $task->status = TaskEntity::STATUS_SUCCESS; $task->lastError = null; $this->saveTask($task); $this->releaseLock($taskId, $lockValue); } catch (Throwable $e) { $task->status = TaskEntity::STATUS_FAILED; $task->lastError = $e->getMessage(); $task->retryCount++; if ($task->retryCount <= $task->maxRetry) { // 重新投递到延迟队列,按指数退避重试 $delay = 5 * (2 ** ($task->retryCount - 1)); $task->nextRunAt = time() + $delay; $this->redis->zAdd('task:pending:zset', $task->nextRunAt, $task->id); $this->saveTask($task); } else { // 进入死信队列 $this->redis->lPush('task:dead:list', $taskId); $this->saveTask($task); } $this->releaseLock($taskId, $lockValue); } } protected function executeTask(TaskEntity $task): mixed { // 根据任务类型分发到具体的业务处理器 // 这里我用一个简单映射做演示,实际项目可用容器或工厂 $handlerClass = 'App\\Tasks\\' . ucfirst($task->type) . 'Task'; if (!class_exists($handlerClass)) { throw new RuntimeException('Task handler not found: ' . $handlerClass); } $handler = new $handlerClass(); return $handler->handle($task->params); } protected function saveTask(TaskEntity $task): void { $this->redis->hSet('task:data:' . $task->id, 'payload', json_encode($task->toArray(), JSON_UNESCAPED_UNICODE)); } protected function releaseLock(string $taskId, string $lockValue): void { // 用 Lua 脚本保证只有持锁者才能删除锁 $script = <<<'LUA' if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end LUA; $this->redis->eval($script, ['task:lock:' . $taskId], [$lockValue]); } }

这段代码里有几个细节很重要:

  1. BRPOP返回的结构是[队列名, 消息],所以取$task[1]才是任务 ID。
  2. 拿锁之后要更新任务状态为EXECUTING,这方便后期排查任务是不是卡住了。
  3. 失败重试我用了指数退避:第一次失败 5 秒后重试,第二次 10 秒,第三次 20 秒。这个策略比固定间隔更适合排障,因为很多任务失败是临时性的,过快重试会加重系统压力。
  4. 死信队列task:dead:list专门收留重试次数耗尽的任务,方便人工介入处理。

3.5 进程管理:用 Supervisor 让 PHP 脚本稳定常驻

PHP CLI 脚本写好后,不可能人工去后台nohup管理,一旦进程崩溃没人拉起来。我用 Supervisor 托管,这是 Linux 下最常用的进程管理工具,能自动重启挂掉的进程,还带日志输出。

在/etc/supervisor/conf.d/task-scheduler.conf写配置:

[program:task-scheduler] command=/usr/bin/php /data/www/app/bin/task-scheduler.php process_name=%(program_name)s_%(process_num)02d numprocs=1 autostart=true autorestart=true redirect_stderr=true stdout_logfile=/data/logs/task-scheduler.log stderr_logfile=/data/logs/task-scheduler.err

消费进程我会多启动几个,让同一台机器上的并发能力更强:

[program:task-worker] command=/usr/bin/php /data/www/app/bin/task-worker.php process_name=%(program_name)s_%(process_num)02d numprocs=4 autostart=true autorestart=true redirect_stderr=true stdout_logfile=/data/logs/task-worker-%(process_num)02d.log stderr_logfile=/data/logs/task-worker.err

如果想在多台服务器上部署消费者,只要把同样的配置复制到另外一台服务器上,它们连接同一个 Redis 即可。这比把 worker 数量直接配置在代码里灵活得多。

4. 深入优化:延迟任务、定时任务、幂等与重复消费

4.1 延迟任务和定时任务到底有什么不同

很多刚接触这块的同学容易把“延迟任务”和“定时任务”混在一起。其实它们在这个架构里的实现方式完全不同:

延迟任务是“一条任务,希望它在某个时间点之后执行一次”。比如下单未支付,30 分钟后关闭。这种任务天然用 ZSet 存,score 是执行时间戳,到期一次执行后任务就完成了。

定时任务是“周期性地执行某个操作”。比如每 5 分钟同步一次缓存,每天凌晨 2 点跑报表。如果只是单机,可以交给 crontab;但如果我们希望它也走统一的失败重试、监控、分布式锁的通道,那就得把 cron 表达式翻译成下一次执行时间,然后像延迟任务一样投递。

我这里给出了一个简化版的实现思路:用cron表达式解析库(比如dragonmantank/cron-expression)知道下一次运行时间,然后调用TaskProducer::push()将任务投递到 ZSet。注意这里投递的“任务类型”需要标记为cron_dispatch,Consumer 在执行这个任务时,会解析对应 cron 配置,并把下一次执行任务再次投递到 ZSet。这样定时任务也被纳入同样的执行和重试体系,而不是散落在 crontab 里了。

// 示例:消费 cron_dispatch 任务时的逻辑 public function handle(array $params): void { $config = $this->findCronConfig($params['cron_key']); $nextRun = CronExpression::factory($config['expr'])->getNextRunDate(); TaskProducer::push('cron_dispatch', [ 'cron_key' => $params['cron_key'], ], $nextRun->getTimestamp() - time()); }

这里又涉及一个难点:如果定时任务配置变更了,旧的cron_dispatch任务已经在 ZSet 里了,要如何取消?实践中我通常给任务加上一个配置版本号,当 Config 表里的 version 变化时,旧任务执行时会发现版本不匹配,直接丢弃。这在业务量不大时足够可靠。

4.2 用 ZSet 实现延迟队列的具体技巧

延迟队列的核心就一句话:入队时间加上延迟时间作为 score,轮询时取出 score 小于等于当前时间的任务。

我把 ZSet 的 key 固定为task:pending:zset,score 是任务的计划执行时间。Scheduler 每 0.1 秒扫一次,把到期的任务搬到 List。如果你有高并发投递需求,可以把这种投递也做成脚本化的原子操作,但我用的是zAdd单条插入,除非单机每秒投递上千条,不然完全没问题。

这里有个性能上的坑:Scheduler 每次循环不能太频繁。我一开始写的是usleep(1000),结果场景里任务不多,调度进程白白占据 Redis 大量连接。后来调成usleep(100000)(100ms),还是能保证任务在 0.1 秒内被扫描到。延迟任务的精度需求一般在秒级甚至分钟级就够了,没必要为了毫秒精度去浪费资源。

4.3 幂等性设计:杜绝重复执行

即使我做了分布式锁,也无法 100% 避免重复执行。比如执行任务过程中网络超时,任务实际已经成功,但 worker 收到了超时异常,代码在 catch 中把任务重新投递到 ZSet,造成二次执行。这种情况在业务上可能是致命的,比如发短信、扣库存、给用户加余额。

所以我会在每个业务处理器的handle方法内,自己再判断一次幂等条件。最典型的做法是在业务表里加一个“任务执行记录”字段,或者使用一张单独的去重表,把任务 ID 作为唯一键。比如:

CREATE TABLE task_execute_log ( id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY, task_id VARCHAR(64) NOT NULL UNIQUE KEY, execute_time DATETIME NOT NULL, result TEXT ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

执行具体业务前先INSERT IGNORE一条记录,影响行数为 0 则说明这个任务已经执行过,直接返回成功。这套方案虽然多一次 DB 写入,但极大地减少了重复执行的风险。如果你是内存敏感型业务,也可以考虑在 Redis 里设置一个task:executed:{taskId}的 key,TTL 设置为业务允许的最长重复时间窗口,但 Redis 数据毕竟不如 DB 那么可信,我更推荐用 DB 做最终幂等。

4.4 失败重试的分级策略

不是所有任务失败都要走指数退避。我把任务分成两档:

  1. 快速重试档。比如同步数据临时网络抖动、外部 API 短暂 5xx。第一次失败后延迟 5 秒,第二次 10 秒,第三次 20 秒,最多重试 3 次。这一类任务要求尽快补上。
  2. 慢速重试档。比如批量导入明细,处理到一半抛异常,马上重试也很可能继续抛异常。这类任务我通常设置延迟 600 秒起,最多重试 5 次,保证不因为高频重试打垮下游。

不同的延迟参数可以在投递时通过TaskEntity的delay和maxRetry指定,完全由业务侧自定义。统一的入口是这个 TaskProducer 和 TaskWorker 的失败流程,而不是业务代码里到处try/catch包裹后自己补任务。

5. 上线实践中的坑与排查技巧

5.1 常见问题速查表

这套系统上线跑了一个多月后,我把遇到的典型问题整理成了一张表,碰到类似情况可以直接对着排查。

现象可能原因排查方法解决方案
任务一直堆积在 ready List,worker 没消费worker 进程挂了ps -ef | grep task-worker检查进程;看 Supervisor 日志自动拉起 worker,检查 Redis 连接是否正常
任务执行了多遍锁过期时间太短,进程实际执行时间超过锁 TTL查看task:lock:对应 key 是否存在,对比任务执行耗时日志按任务最长耗时调整锁 TTL,或者用 Redisson 类似的看门狗机制
延迟任务到了时间却不执行ZSet 中 score 未正确更新在 Scheduler 脚本里打日志,确认zRangeByScore拉到的数量检查投递时time() + delay是否正确,Scheduler 是否被锁卡住
某个任务重试次数耗光后被丢了正常进入死信队列lLen task:dead:list查看积压数量增加后台管理脚本,消费死信队列并人工修复
Redis 内存飙升未清理已完成的任务查看task:data:*key 总数定期扫描并清理 N 天前成功的任务数据
worker 大量 500 报错业务逻辑 bug,导致任务疯狂重试观察日志中 Exception 堆栈修复 bug,并考虑加入熔断机制

5.2 如何优雅处理 PHP 内存泄漏和长驻进程资源释放

PHP 写 long-running 的 CLI 进程,最担心的就是内存泄漏。我们这套任务调度里,如果TaskEntity::fromArray()或业务处理器内部,每次循环都无意间引用了新的全局对象而忘记释放,内存会慢慢上涨。我用了一个很土但有效的办法:

在run()循环体的末尾,每隔一定循环次数主动调用gc_collect_cycles()和检查memory_get_usage()。如果内存超过基准的 120%,就在处理完当前任务后exit(0),让 Supervisor 自动重新拉起一个干净进程。

$baseMemory = memory_get_usage(); $counter = 0; while (true) { // ... 处理任务 if (++$counter % 100 === 0) { gc_collect_cycles(); $currentMemory = memory_get_usage(); if ($currentMemory - $baseMemory > 100 * 1024 * 1024) { // 内存增长超过 100MB,主动退出,由 Supervisor 重启 log("worker restart due to memory limit: {$currentMemory}"); exit(0); } } }

这个方法看起来笨,但在 PHP 生态里非常实用。尤其是你使用了很多第三方 SDK,它们可能在某些调用链上产生静态缓存,不主动清理就会一直膨胀。

5.3 监控与日志:让调度器可观测

任务调度系统如果不可观测,线上出了问题是灾难。我加的监控非常简单,但都是有效果的:

  • llen task:ready:list入监控,超过阈值报警,比如持续 5 分钟大于 5000,说明消费能力不足。
  • llen task:dead:list入监控,大于 0 就报警并拉出前几条,人工介入。
  • 每执行完一条任务,在日志里记录任务 ID、类型、耗时、状态,按天分文件。
  • 使用 Redis 的SCAN定期统计task:data:*总量和类型分布,方便定位某类任务是否异常偏多。

不要小看这些基础指标,有了它们之后,我再也不用半夜爬起来翻日志找问题了。

5.4 扩展与后续思考

这套方案目前支撑了我们每天几十万量级的定时任务,Redis 峰值内存 200MB 不到,worker 节点加到 6 个进程后,消费延迟控制在一秒以内。如果你的业务量比我更大,需要再往两个方向升级:一是用 Redis Cluster 或 Codis 做存储分片,把队列分散到多组 Redis;二是把任务调度数据从 Redis 持久化到数据库,定期归档,减少 Redis 内数据量。

还有一点,很多业务方希望“某个任务已经进入队列但还没到时候,就想立即执行”,我给 TaskProducer 加了一个triggerNow方法,把任务的 score 改成当前时间,同时更新 Hash 里的计划执行时间。这样运维同学不需要重新投递,只需要在后台点一下“立即执行”按钮,非常实用。

踩过的坑多了之后,我最大的感受是:不要把任务调度做成一个黑盒。它应该是团队里每个后端都能看懂、能修改、能排查的基础组件。这套基于 Redis 的 PHP 分布式任务调度,最大的价值不是它的性能上限有多高,而是它足够简单,让每个接手的人都能在半小时内读懂全部代码。后来我们写代码的时候有个约定:凡是任何一种任务出现异常,第一步先看任务的完整生命周期状态,第二步看task:lock有没有残留,第三步看死信队列。靠这三个动作,几乎能解决 90% 的调度问题。

如果你也在评估 PHP 分布式调度的方案,可以照着我这套思路先搭一个最小版本跑跑看,在真实业务压力下观察一段时间再逐步完善。分布式调度最怕上来就追求大而全,先把核心链路跑通,再一点点加细节,反而走得更稳。

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

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

立即咨询