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

CodeIgniter中RabbitMQ多文件上传并行处理优化问题咨询

问题描述

我在CodeIgniter项目中集成RabbitMQ搭建了多文件上传系统,支持用户批量上传文件。文件上传完成后,其ID会被加入RabbitMQ队列,由后台执行OCR扫描等任务。通过终端执行命令php index.php JobProcessing/process_progress_rabbitmq启动worker处理任务时,单用户场景运行正常,但当多用户(如50人同时上传)操作时,每个用户的任务需等待前序任务完成,无法并行处理。

我考虑通过创建多个worker解决此问题,但担心worker数量存在上限,比如当有500用户同时上传时,可安全创建多少个worker?同时想了解结合系统资源与RabbitMQ限制的推荐处理方案。

实现代码
<?php

defined('BASEPATH') or exit('No direct script access allowed');
require_once(APPPATH . 'third_party/Rabbit_mq/vendor/autoload.php');
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;
use Spatie\Async\Pool;

class Rabbit_mq
{
    protected $ci, $connection, $channel, $pool;

    public function __construct()
    {
        try {
            $this->ci = &get_instance();
            $this->ci->load->model('invoice_scan_model');
            $this->ci->load->library('scan_invoice_lib');
            $this->ci->load->library('quick_books');
            $this->connection = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
            $this->channel = $this->connection->channel();
        } catch (Exception $e) {
            echo "RabbitMQ Connection Error: " . $e->getMessage() . "\n";
        }
    }

    function addToQueue($data)
    {
        $this->channel->queue_declare('file_processing_new', false, true, false, false);

        $msg = new AMQPMessage(json_encode($data));
        $this->channel->basic_publish($msg, '', 'file_processing_new');
    }

    public function processQueue()
    {
        $this->channel->queue_declare('file_processing_new', false, true, false, false);

        // Set prefetch count to allow multiple messages to be handled concurrently
        $this->channel->basic_qos(null, 15, null);

        $callback = function ($msg) {
            $data = json_decode($msg->body, true);
            $fileId = $data['file_id'];
            $filePath = $data['file_path'];

            try {
                echo "Received message for file ID: $fileId\n";

                // Initialize the async pool
                $this->pool = Pool::create();

                // Add task to the pool for parallel processing
                $this->pool->add(function () use ($fileId, $filePath) {
                    try {
                        echo "Processing file ID: $fileId\n";
                        $data = array(
                            'progress_status' => 'Processing',
                            'progress_percentage' => 50,
                        );
                        $this->ci->invoice_scan_model->update_data(['id' => $fileId], 'tblapi_save_invoice_file', $data);
                        echo "File ID: $fileId marked as Processing.\n";

                        // Simulate file processing (OCR or other logic here)
                        $this->process_invoice($fileId);

                        // Update status to 'Completed' after processing
                        $data = array(
                            'progress_status' => 'Completed',
                            'progress_percentage' => 100,
                        );
                        $this->ci->invoice_scan_model->update_data(['id' => $fileId], 'tblapi_save_invoice_file', $data);
                        echo "File ID: $fileId processing completed.\n";
                    } catch (Exception $e) {
                        echo "Error processing file ID: $fileId - " . $e->getMessage() . "\n";
                    }
                });

                // Wait for all tasks to finish
                $this->pool->wait(); // This will wait until all tasks are done

                // Acknowledge the message after all tasks are processed
                $this->channel->basic_ack($msg->delivery_info['delivery_tag']);
                echo "Acknowledging message for file ID: $fileId\n";
            } catch (Exception $e) {
                echo "Error processing file ID: $fileId - " . $e->getMessage() . "\n";
                $this->channel->basic_nack($msg->delivery_info['delivery_tag']);
            }
        };

        // Multiple consumers (workers) consuming the messages
        for ($i = 0; $i < 25; $i++) {
            $this->channel->basic_consume('file_processing_new', '', false, false, false, false, $callback);
        }

        echo "Waiting for messages. To exit press CTRL+C\n";

        // Consume messages concurrently by multiple workers
        while ($this->channel->callbacks) {
            $this->channel->wait();
        }

        // Close the channel and connection when done
        $this->channel->close();
        $this->connection->close();
    }
}
问题分析与解决方案

1. 当前无法并行的原因

你的代码存在两个核心问题导致任务串行执行:

  • 单进程消费者阻塞:虽然循环创建了25个消费者回调,但所有回调都运行在同一个PHP进程中,PHP的channel->wait()是单线程阻塞逻辑,同一时间只能处理一个消息回调,其他消息只能排队等待。
  • 异步池误用:每个消息回调里创建独立异步池,仅添加一个任务后立刻调用pool->wait()等待任务完成,相当于把同步任务包装了一层异步壳,本质还是阻塞当前消息的处理流程,没有实现真正的并行。

2. Worker数量的安全上限

Worker数量没有固定值,需结合服务器资源和RabbitMQ配置判断:

  • CPU核心数:OCR属于CPU密集型任务,worker数量建议不超过CPU核心数的1.52倍(比如8核服务器设置1216个worker)。超过这个范围会导致进程频繁上下文切换,反而降低处理效率。
  • 内存资源:每个PHP worker进程通常占用几十到上百MB内存,要确保服务器剩余内存能支撑所有worker运行,避免出现内存不足(OOM)。
  • RabbitMQ限制:RabbitMQ默认连接数上限是65535,每个独立worker进程会占用一个连接,这个限制一般不会触发,但要注意RabbitMQ的内存、文件句柄配置,避免因连接过多导致服务崩溃。

针对500用户同时上传的场景,不需要创建500个worker。比如8核服务器设置10~15个worker即可,RabbitMQ会自动将队列消息均匀分发给各个worker并行处理,未处理的消息会在队列中积压,worker会持续消费直到队列清空。

3. 推荐处理方案

方案一:多独立Worker进程(最可靠)

放弃单进程内多消费者的逻辑,改为启动多个独立的worker进程:

  • 每次执行php index.php JobProcessing/process_progress_rabbitmq启动一个worker进程,手动启动多个(比如10个),或者用Supervisor等进程管理工具自动维护固定数量的worker。
  • 修改processQueue方法,移除创建25个消费者的循环,每个worker只创建一个消费者:
    // 移除原for循环,仅保留单个消费者
    $this->channel->basic_consume('file_processing_new', '', false, false, false, false, $callback);
    
  • 保留basic_qos配置,设置合理的prefetch count(比如每个worker同时处理5~10条消息,根据任务耗时调整),避免单个worker一次性获取过多消息导致积压。

方案二:优化异步任务逻辑(单进程多线程)

如果坚持用单进程处理,需优化Spatie Async的使用方式:

  • 在processQueue方法初始化全局异步池,设置最大并发数(等于CPU核心数)。
  • 消息回调中直接将任务添加到全局池,不调用wait(),而是在消息循环等待过程中定期检查任务状态。不过这种方式在PHP中实现复杂度高,可靠性不如多进程方案。

方案三:队列拆分与分组调度

若需要精细调度,可将主队列拆分为多个子队列(比如按文件类型、用户分组),启动不同的worker组消费对应队列,避免某类任务占用所有worker资源。

额外优化建议

  • 消息确认机制:确保任务真正完成后再调用basic_ack,避免任务失败导致消息丢失;失败任务可通过basic_nack重新放回队列,或配置死信队列专门处理失败消息。
  • 监控与日志:添加详细日志记录,监控worker运行状态和队列长度,及时调整worker数量。
  • 资源限制:给每个worker进程设置CPU、内存限制,避免单个进程占用过多资源影响其他worker运行。

内容的提问来源于stack exchange,提问作者vinod

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 17:04:55