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
相关产品推荐
相关产品推荐

