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

Magento批量处理美国税率记录时curl_multi请求异常求助

问题描述

需要在Magento中为美国所有州、城市、邮编组合设置税率,总计13万条记录,每年执行两次。目前采用RabbitMQ实现方案:

  • Cron Job按500条批量发布消息至消费者
  • 消费者接收消息后,使用curl_multi_init请求指定URL,根据响应结果更新税率数据表

遇到的异常:

  • 当记录数≤8条时,curl请求正常,可获取状态码及完整响应
  • 当记录数>10条时,响应不稳定,偶尔返回空值
  • 当记录数>15条时,响应始终为空,curl状态码持续为0

已尝试方案:使用标准PHP代码处理,问题依旧存在。

消费者代码

private function getTransactionCallback(QueueInterface $queue)
{
    return function (EnvelopeInterface $message) use ($queue) {
        /** @var LockInterface $lock */
        $lock = null;
        try {
            $lock = $this->messageController->lock($message, $this->configuration->getConsumerName());
            $body = $message->getBody();
            
            $writer = new \Zend_Log_Writer_Stream(BP . '/var/log/consumer.log');
            $logger = new \Zend_Log();
            $logger->addWriter($writer);
    
            try{
                $logger->info('Start');
                $start = microtime(true);
                $body = json_decode(json_decode($body, true),true);
                $mh = curl_multi_init();
                $handles = [];
                $associatedData = [];
                
                foreach ($body as $key => $msg) {
                    $ch = curl_init();
                    $handles[] = $ch;
                    curl_setopt($ch, CURLOPT_URL,$msg['endpoints']);
                    curl_setopt($ch, CURLOPT_RETURNTRANSFER, 1);
                    curl_setopt($ch, CURLOPT_TIMEOUT, 5);

                    $associatedData[(int)$ch] = $msg['entity_id'];

                    curl_multi_add_handle($mh,$ch);

                }

                $running = null;
                do {
                    curl_multi_exec($mh, $running);
                } while ($running);

                $objectManager = \Magento\Framework\App\ObjectManager::getInstance();
                $sql = $objectManager->create('Magento\Framework\App\ResourceConnection');
                $connection = $sql->getConnection();
                $tableName = $sql->getTableName('us_tax');
                foreach($handles as $ch){
                    $result = curl_multi_getcontent($ch);
                    $rate = 0;
                   
                    try {
                        $id = (int)$ch;
                        $status = curl_getinfo($ch, CURLINFO_HTTP_CODE);
                        $logger->info('ID : '.$associatedData[$id]);
                        $logger->info('Rate : '.$result);
                        $logger->info(print_r($status,true));
                        if( $result ){
                            $dd = json_decode($result,true);
                            $rate = ((float)json_decode(json_decode($dd))) * 100;     
                        }

                        $connection->update($tableName,['tax_rate' => $rate],['entity_id = ?' => $associatedData[$id]]);

                    }catch(\Exception $e){
                        $logger->info('Error : '.$e->getMessage());
                    }

                    curl_multi_remove_handle($mh, $ch);

                }

                curl_multi_close($mh);
                $end = microtime(true);
                $time = number_format(($end - $start), 2);
                $logger->info('Time execute : '.$time);
                $logger->info('End');

            }catch (\Exception $e){
                $logger->info('Error : '.$e->getMessage());

            }

            $data = true;
            if ($data === false) {
                $queue->reject($message); // if get error in message process
            }
            $queue->acknowledge($message); // send acknowledge to queue 
        } catch (MessageLockException $exception) {
            $queue->acknowledge($message);
        } catch (ConnectionLostException $e) {
            $queue->acknowledge($message);
            if ($lock) {
                $this->resource->getConnection()
                    ->delete($this->resource->getTableName('queue_lock'), ['id = ?' => $lock->getId()]);
            }
        } catch (NotFoundException $e) {
            $queue->acknowledge($message);
            $this->logger->warning($e->getMessage());
        } catch (\Exception $e) {
            $queue->reject($message, false, $e->getMessage());
            $queue->acknowledge($message);
            if ($lock) {
                $this->resource->getConnection()
                    ->delete($this->resource->getTableName('queue_lock'), ['id = ?' => $lock->getId()]);
            }
        }
    };
}

解决方案建议

一、修复curl_multi请求异常

1. 限制并发请求数

一次性添加过多handle会触发目标服务器限流或PHP资源耗尽,建议将并发数控制在10以内,分批处理请求:

// 替换原批量添加handle的逻辑,改为分批处理
$batchSize = 10;
$totalItems = count($body);

for ($i = 0; $i < $totalItems; $i += $batchSize) {
    $batch = array_slice($body, $i, $batchSize);
    $mh = curl_multi_init();
    $handles = [];
    $associatedData = [];

    foreach ($batch as $msg) {
        $ch = curl_init();
        curl_setopt($ch, CURLOPT_URL, $msg['endpoints']);
        curl_setopt($ch, CURLOPT_RETURNTRANSFER, 1);
        curl_setopt($ch, CURLOPT_TIMEOUT, 10); // 延长超时时间
        curl_setopt($ch, CURLOPT_CONNECTTIMEOUT, 5); // 添加连接超时设置
        curl_setopt($ch, CURLOPT_FOLLOWLOCATION, true); // 自动处理重定向
        
        $associatedData[(int)$ch] = $msg['entity_id'];
        curl_multi_add_handle($mh, $ch);
        $handles[] = $ch;
    }

    // 优化curl_multi_exec循环,减少CPU占用
    $running = null;
    do {
        $status = curl_multi_exec($mh, $running);
        if ($status !== CURLM_OK) break;
        curl_multi_select($mh); // 等待活动连接,避免空循环
    } while ($running > 0);

    // 处理响应
    foreach ($handles as $ch) {
        $result = curl_multi_getcontent($ch);
        $rate = 0;
        try {
            $id = (int)$ch;
            $statusCode = curl_getinfo($ch, CURLINFO_HTTP_CODE);
            $curlError = curl_error($ch);
            
            $logger->info('ID : ' . $associatedData[$id]);
            $logger->info('HTTP Status : ' . $statusCode);
            if ($curlError) $logger->info('Curl Error : ' . $curlError);

            // 仅处理200状态且无错误的响应
            if ($result && $statusCode == 200 && empty($curlError)) {
                $dd = json_decode($result, true);
                // 根据实际接口返回格式调整解析逻辑,避免多层json_decode嵌套
                $rate = ((float)$dd) * 100;
            }

            $connection->update($tableName, ['tax_rate' => $rate], ['entity_id = ?' => $associatedData[$id]]);
        } catch (\Exception $e) {
            $logger->info('Processing ID ' . $associatedData[(int)$ch] . ' failed: ' . $e->getMessage());
        }

        curl_multi_remove_handle($mh, $ch);
        curl_close($ch); // 显式关闭单个handle释放资源
    }

    curl_multi_close($mh);
}

2. 完善错误日志

添加curl错误信息捕获,通过curl_error($ch)获取具体失败原因(如连接超时、服务器拒绝等),便于定位问题。

3. 调整服务器配置

  • 检查PHP的max_execution_time,确保脚本有足够时间处理批量请求
  • 若目标服务器有频率限制,可在请求中添加随机延迟(如sleep(0.1))避免触发限流

二、替代方案建议

1. 使用Magento原生HTTP客户端

替换手动curl_multi实现,改用Magento自带的Magento\Framework\HTTP\Client\Curl或异步HTTP客户端(Magento 2.4+支持),利用框架的连接池和资源管理能力,降低手动维护成本。

2. 缩小任务粒度

将500条的批量消息拆分为10条/批,减少单个消费者的处理压力,同时可通过增加消费者进程实现横向扩展,提升整体处理效率。

3. 离线预处理数据

由于任务每年仅执行两次,可提前从税率接口导出全量美国地区税率数据,生成CSV文件后通过Magento批量导入工具直接更新数据表,彻底避免实时请求接口的不稳定问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 03:04:57