You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.02 22:32:32