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

如何在Magento 2后台为商品CSV导入实现队列功能

Magento 2 实现CSV商品导入队列化处理

要实现后台CSV商品导入时的队列功能,核心思路是拦截默认的导入验证/执行流程,将导入任务放入消息队列,让后台消费者异步执行导入操作,前端仅返回队列任务提交成功的提示。以下是具体实现步骤:

1. 创建自定义模块

先搭建基础模块结构,命名为Vendor_ImportQueue:

app/code/Vendor/ImportQueue/registration.php

<?php
use Magento\Framework\Component\ComponentRegistrar;

ComponentRegistrar::register(ComponentRegistrar::MODULE, 'Vendor_ImportQueue', __DIR__);

app/code/Vendor/ImportQueue/etc/module.xml

<?xml version="1.0"?>
<config xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="urn:magento:framework:Module/etc/module.xsd">
    <module name="Vendor_ImportQueue" setup_version="1.0.0">
        <sequence>
            <module name="Magento_ImportExport"/>
        </sequence>
    </module>
</config>

2. 重写导入控制器

默认商品导入流程会在验证后直接执行导入,我们需要重写Validate控制器,将任务推入队列:

app/code/Vendor/ImportQueue/etc/adminhtml/di.xml

<?xml version="1.0"?>
<config xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="urn:magento:framework:ObjectManager/etc/config.xsd">
    <preference for="Magento\ImportExport\Controller\Adminhtml\Import\Validate" type="Vendor\ImportQueue\Controller\Adminhtml\Import\Validate"/>
</config>

app/code/Vendor/ImportQueue/Controller/Adminhtml/Import/Validate.php

<?php
namespace Vendor\ImportQueue\Controller\Adminhtml\Import;

use Magento\Framework\App\Action\HttpPostActionInterface;
use Magento\ImportExport\Controller\Adminhtml\Import\Validate as OriginalValidate;
use Magento\Framework\Queue\QueueInterface;
use Magento\Framework\Serialize\SerializerInterface;

class Validate extends OriginalValidate implements HttpPostActionInterface
{
    protected $queue;
    protected $serializer;

    public function __construct(
        \Magento\Backend\App\Action\Context $context,
        \Magento\ImportExport\Model\Import $importModel,
        \Magento\ImportExport\Model\Import\Factory $importFactory,
        \Magento\ImportExport\Model\Export\Factory $exportFactory,
        \Magento\ImportExport\Model\Import\Config $importConfig,
        \Magento\Framework\App\Config\ScopeConfigInterface $scopeConfig,
        QueueInterface $queue,
        SerializerInterface $serializer
    ) {
        parent::__construct($context, $importModel, $importFactory, $exportFactory, $importConfig, $scopeConfig);
        $this->queue = $queue;
        $this->serializer = $serializer;
    }

    public function execute()
    {
        $this->_initSession();
        $data = $this->getRequest()->getPostValue();
        
        // 执行默认验证逻辑
        $validationResult = $this->_importModel->validateSource();
        
        if (!$validationResult) {
            // 验证失败,返回原错误提示
            return parent::execute();
        }

        // 验证通过,准备队列任务数据
        $importData = [
            'entity' => $data['entity'],
            'file' => $this->_importModel->getUploadedFileName(),
            'behavior' => $data['behavior'],
            'validation_strategy' => $data['validation_strategy'],
            'allowed_error_count' => $data['allowed_error_count'],
            'import_images_file_dir' => $data['import_images_file_dir'] ?? ''
        ];

        // 将任务推入队列
        $this->queue->push($this->serializer->serialize($importData));

        // 返回队列提交成功提示
        $this->messageManager->addSuccessMessage(__('商品导入任务已加入队列,将在后台自动执行,请稍后查看导入结果。'));
        
        return $this->_redirect('importexport/import/index');
    }
}

3. 配置消息队列

定义队列名称和消费者,让Magento识别我们的队列任务:

app/code/Vendor/ImportQueue/etc/queue.xml

<?xml version="1.0"?>
<config xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="urn:magento:framework-message-queue:etc/queue.xsd">
    <queue name="import_product_csv" consumer="import_product_csv_consumer"/>
</config>

app/code/Vendor/ImportQueue/etc/crontab.xml(可选)

如果需要定时触发消费者,添加crontab配置:

<?xml version="1.0"?>
<config xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="urn:magento:module:Magento_Cron:etc/crontab.xsd">
    <group id="default">
        <job name="import_product_csv_consumer" instance="Vendor\ImportQueue\Model\Consumer" method="process">
            <schedule>* * * * *</schedule>
        </job>
    </group>
</config>

4. 创建消费者类处理后台导入

消费者类负责从队列中取出任务,执行实际的商品导入逻辑:

app/code/Vendor/ImportQueue/Model/Consumer.php

<?php
namespace Vendor\ImportQueue\Model;

use Magento\Framework\Queue\ConsumerInterface;
use Magento\Framework\Serialize\SerializerInterface;
use Magento\ImportExport\Model\ImportFactory;
use Magento\Framework\App\Filesystem\DirectoryList;
use Magento\Framework\Filesystem;
use Psr\Log\LoggerInterface;

class Consumer implements ConsumerInterface
{
    protected $serializer;
    protected $importFactory;
    protected $filesystem;
    protected $logger;

    public function __construct(
        SerializerInterface $serializer,
        ImportFactory $importFactory,
        Filesystem $filesystem,
        LoggerInterface $logger
    ) {
        $this->serializer = $serializer;
        $this->importFactory = $importFactory;
        $this->filesystem = $filesystem;
        $this->logger = $logger;
    }

    public function process($message)
    {
        try {
            $importData = $this->serializer->unserialize($message);
            
            // 获取临时上传目录
            $tmpDir = $this->filesystem->getDirectoryWrite(DirectoryList::VAR_DIR)->getAbsolutePath('importexport/');
            $filePath = $tmpDir . $importData['file'];

            // 初始化导入模型
            $importModel = $this->importFactory->create();
            $importModel->setData([
                'entity' => $importData['entity'],
                'behavior' => $importData['behavior'],
                'validation_strategy' => $importData['validation_strategy'],
                'allowed_error_count' => $importData['allowed_error_count'],
                'import_images_file_dir' => $importData['import_images_file_dir'],
                'source' => $filePath
            ]);

            // 执行导入
            $importModel->importSource();
            
            // 记录成功日志
            $this->logger->info('CSV商品导入任务执行成功,文件:' . $importData['file']);
        } catch (\Exception $e) {
            // 记录错误日志
            $this->logger->error('CSV商品导入任务执行失败:' . $e->getMessage());
        }
    }
}

5. 配置消费者关系

将消费者类与队列绑定:

app/code/Vendor/ImportQueue/etc/di.xml

<?xml version="1.0"?>
<config xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="urn:magento:framework:ObjectManager/etc/config.xsd">
    <type name="Magento\Framework\MessageQueue\ConsumerInterface">
        <arguments>
            <argument name="consumers" xsi:type="array">
                <item name="import_product_csv_consumer" xsi:type="array">
                    <item name="connection" xsi:type="string">db</item>
                    <item name="queue" xsi:type="string">import_product_csv</item>
                    <item name="consumerInstance" xsi:type="object">Vendor\ImportQueue\Model\Consumer</item>
                    <item name="maxMessages" xsi:type="number">10</item>
                    <item name="batchSize" xsi:type="number">1</item>
                </item>
            </argument>
        </arguments>
    </type>
</config>

6. 模块启用与测试

  1. 执行命令启用模块:
bin/magento module:enable Vendor_ImportQueue
bin/magento setup:upgrade
bin/magento setup:di:compile
bin/magento cache:flush
  1. 后台上传CSV文件并提交导入,此时会收到队列任务成功的提示,导入操作在后台由消费者执行。
  2. 可通过var/log/system.log查看导入日志,确认任务执行状态。

内容的提问来源于stack exchange,提问作者abhishek singh developer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 07:16:11