Beanstalkd队列Worker意外冻结需重启,请求故障排查帮助
Beanstalkd PHP Worker无报错静默停止的排查与修复方案
问题背景
部署3个独立Beanstalkd队列,每个队列对应无框架PHP Worker,通过Pheanstalk库交互。近期Worker会无报错突然停止,重启后可正常运行数日,随后重复该循环。
原始Worker代码
<?php ini_set('display_errors', 1); ini_set('display_startup_errors', 1); error_reporting(E_ALL); require 'vendor/autoload.php'; use Pheanstalk\Pheanstalk; function function_name() { $log_file = './log/error.log'; $tube_name = 'current_tube_name'; $memoryLimit = 128; $dotenv = Dotenv\Dotenv::createImmutable(__DIR__); $dotenv->load(); $beanstalkd_host = $_ENV['BEANSTALKD_HOST']; $beanstalkd_port = $_ENV['BEANSTALKD_PORT']; while (true) { try { $pheanstalk = Pheanstalk::create($beanstalkd_host); $pheanstalk->watch($tube_name)->ignore('default'); } catch(\Exception $exception) { $error_message = "CRITICAL|" . __FILE__ . "|" . __LINE__ . "|" . "Failed to initiate queue worker; could not connect to beanstalkd.Error {$exception->getMessage()}"; error_log($error_message, 3, $log_file); exit; } $job = $pheanstalk->reserve(10); if (!$job) { sleep(5); continue; //move on to next iteration } $job_in_queue = $job->getData();//outputting the message $job_pay_load = json_decode($job_in_queue, true); if(is_array($job_pay_load) == false) { $job_pay_load = json_decode(unserialize($job_in_queue), true); } if (is_null($job_pay_load) || !is_array($job_pay_load) || count($job_pay_load) === 0 || !is_countable($job_pay_load)) { $error_message = "CRITICAL|" . __FILE__ . "|" . __LINE__ . "|" . "Job payload structure is not correct. Data: " . print_r($jobPayload, true); error_log($error_message, 3, $log_file); $pheanstalk->bury($job); $job = null; $job_pay_load = null; continue; } //do the actual process $payload = array( 'key1' => $val1, 'key2' => $val2 ); $next_tube_name = 'second_tube'; $json_payload = json_encode($payload); $pheanstalk ->useTube($next_tube_name) ->put($json_payload); $pheanstalk->delete($job); $job = null; $json_payload = null; $job_pay_load = []; $job_in_queue = null; gc_collect_cycles(); if((memory_get_usage() / 1024 / 1024) >= $memoryLimit) { $error_message = "We ran out of memory, restarting... in file " . __FILE__ . " on line" . __LINE__; error_log($error_message, 3, $log_file); exit; } } } function_name();
核心排查方向
1. 连接资源泄漏
- 原始代码每次循环都新建Pheanstalk连接,未关闭旧连接,长时间运行会耗尽系统文件句柄,导致Worker无法读写数据静默退出。
- 修复:将连接实例化移到循环外,仅在连接失效时重建。
2. 未捕获的致命错误
- 虽然开启了错误报告,但
reserve()或业务逻辑中可能存在无法被try/catch捕获的致命错误(如PHP扩展崩溃、内存耗尽在检查前触发),导致Worker直接退出无日志。 - 优化:添加全局错误/ shutdown捕获,记录所有致命错误。
3. 内存泄漏
- 手动GC和内存检查存在漏洞,业务逻辑中未释放的大对象/循环引用会导致内存缓慢泄漏,最终被系统OOM Killer杀死(无PHP层面日志)。
- 优化:使用
memory_get_usage(true)获取实际分配内存,将内存检查移到业务逻辑执行后,确保大变量使用后及时置空。
4. Beanstalkd连接静默断开
reserve(10)超时后,若Beanstalkd主动断开连接,PHP可能不抛出异常,后续连接重建失败会触发exit退出。- 修复:每次循环前检查连接有效性,失效则重建。
5. 系统层面限制
- 检查系统日志(
/var/log/syslog或/var/log/messages)是否有OOM Killer杀死Worker的记录(关键词oom-killer)。 - 执行
ulimit -n检查文件句柄限制,确保Worker有足够句柄数。
优化后的Worker代码
<?php ini_set('display_errors', 1); ini_set('display_startup_errors', 1); error_reporting(E_ALL); require 'vendor/autoload.php'; use Pheanstalk\Pheanstalk; function function_name() { $log_file = './log/error.log'; $tube_name = 'current_tube_name'; $memoryLimit = 128; // MB $dotenv = Dotenv\Dotenv::createImmutable(__DIR__); $dotenv->load(); $beanstalkd_host = $_ENV['BEANSTALKD_HOST']; $beanstalkd_port = $_ENV['BEANSTALKD_PORT']; // 全局错误捕获,记录所有非致命错误 set_error_handler(function($errno, $errstr, $errfile, $errline) use ($log_file) { $error_message = "ERROR|" . $errfile . "|" . $errline . "|" . $errstr; error_log($error_message, 3, $log_file); // 致命错误直接退出 if (in_array($errno, [E_ERROR, E_PARSE, E_CORE_ERROR, E_COMPILE_ERROR])) { exit(1); } }); // 捕获脚本终止时的致命错误 register_shutdown_function(function() use ($log_file) { $error = error_get_last(); if ($error && in_array($error['type'], [E_ERROR, E_PARSE, E_CORE_ERROR, E_COMPILE_ERROR])) { $error_message = "FATAL|" . $error['file'] . "|" . $error['line'] . "|" . $error['message']; error_log($error_message, 3, $log_file); } }); $pheanstalk = null; // 连接初始化/重建函数 $initConnection = function() use (&$pheanstalk, $beanstalkd_host, $tube_name, $log_file) { try { $pheanstalk = Pheanstalk::create($beanstalkd_host); $pheanstalk->watch($tube_name)->ignore('default'); return true; } catch(\Exception $exception) { $error_message = "CRITICAL|" . __FILE__ . "|" . __LINE__ . "|" . "Failed to connect to beanstalkd: {$exception->getMessage()}"; error_log($error_message, 3, $log_file); return false; } }; // 初始化连接失败直接退出 if (!$initConnection()) { exit(1); } while (true) { // 检查连接有效性,失效则重建 if (!$pheanstalk || !$pheanstalk->getConnection()->isConnected()) { if (!$initConnection()) { sleep(5); continue; } } try { $job = $pheanstalk->reserve(10); } catch(\Exception $e) { $error_message = "CRITICAL|" . __FILE__ . "|" . __LINE__ . "|" . "Reserve job failed: {$e->getMessage()}"; error_log($error_message, 3, $log_file); $pheanstalk = null; // 标记连接失效 sleep(5); continue; } if (!$job) { sleep(5); continue; } $job_in_queue = $job->getData(); $job_pay_load = json_decode($job_in_queue, true); if (!is_array($job_pay_load)) { $job_pay_load = json_decode(unserialize($job_in_queue), true); } // 修复原代码变量名错误:$jobPayload改为$job_in_queue if (is_null($job_pay_load) || !is_array($job_pay_load) || empty($job_pay_load)) { $error_message = "CRITICAL|" . __FILE__ . "|" . __LINE__ . "|" . "Invalid job payload: " . print_r($job_in_queue, true); error_log($error_message, 3, $log_file); $pheanstalk->bury($job); // 及时清理变量 $job = null; $job_pay_load = null; $job_in_queue = null; continue; } // 业务逻辑执行(替换为实际代码) $val1 = $job_pay_load['key1'] ?? ''; $val2 = $job_pay_load['key2'] ?? ''; $payload = array( 'key1' => $val1, 'key2' => $val2 ); $next_tube_name = 'second_tube'; $json_payload = json_encode($payload); try { $pheanstalk ->useTube($next_tube_name) ->put($json_payload); } catch(\Exception $e) { $error_message = "CRITICAL|" . __FILE__ . "|" . __LINE__ . "|" . "Put job to {$next_tube_name} failed: {$e->getMessage()}"; error_log($error_message, 3, $log_file); $pheanstalk->bury($job); $job = null; continue; } try { $pheanstalk->delete($job); } catch(\Exception $e) { $error_message = "CRITICAL|" . __FILE__ . "|" . __LINE__ . "|" . "Delete job failed: {$e->getMessage()}"; error_log($error_message, 3, $log_file); } // 清理所有临时变量 $job = null; $json_payload = null; $job_pay_load = []; $job_in_queue = null; $payload = null; $val1 = null; $val2 = null; gc_collect_cycles(); // 使用实际分配的内存判断,更准确 if ((memory_get_usage(true) / 1024 / 1024) >= $memoryLimit) { $error_message = "MEMORY_LIMIT|" . __FILE__ . "|" . __LINE__ . "|" . "Memory limit reached, restarting..."; error_log($error_message, 3, $log_file); exit(1); } } } function_name();
内容的提问来源于stack exchange,提问作者Krish
相关产品推荐
相关产品推荐

