定制 iq2i/messenger-import-bundle 二次开发

按需修改功能、优化性能、对接业务系统,提供一站式技术支持

邮箱:yvsm@zunyunkeji.com | QQ:316430983 | 微信:yvsm316

iq2i/messenger-import-bundle

Composer 安装命令:

composer require iq2i/messenger-import-bundle

包简介

A Symfony bundle that tracks the completion of asynchronous imports dispatched via symfony/messenger

README 文档

README

A Symfony bundle that tracks the completion of asynchronous imports dispatched via symfony/messenger.

When importing large files, dispatching messages asynchronously speeds up processing but makes it impossible to know when all messages have been handled. This bundle solves that problem by listening to Messenger worker events and detecting when every message in a batch has been processed — whether successfully or not.

Requirements

  • PHP 8.1+
  • Symfony 6.4 / 7.x / 8.x
  • symfony/messenger
  • Doctrine ORM (for the provided traits)

Installation

composer require iq2i/messenger-import-bundle

Register the bundle in config/bundles.php:

return [
    // ...
    IQ2i\MessengerImportBundle\MessengerImportBundle::class => ['all' => true],
];

How it works

  1. Before dispatching messages, you create an ImportBatch entity that stores the total number of messages to process.
  2. Each message carries the batch ID via BatchAwareMessageInterface.
  3. The bundle's subscriber listens to WorkerMessageHandledEvent and WorkerMessageFailedEvent. After each message, it decrements the batch counter.
  4. When the counter reaches zero, an ImportBatchCompletedEvent is dispatched. You listen to this event to send a notification, trigger a follow-up action, etc.

Setup

1. Create the batch entity

Create an entity that implements ImportBatchInterface and uses ImportBatchTrait:

// src/Entity/ImportBatch.php

use Doctrine\ORM\Mapping as ORM;
use IQ2i\MessengerImportBundle\Model\ImportBatchInterface;
use IQ2i\MessengerImportBundle\Model\ImportBatchTrait;

#[ORM\Entity(repositoryClass: ImportBatchRepository::class)]
class ImportBatch implements ImportBatchInterface
{
    use ImportBatchTrait;

    #[ORM\Id]
    #[ORM\GeneratedValue(strategy: 'CUSTOM')]
    #[ORM\CustomIdGenerator(class: UuidGenerator::class)]
    #[ORM\Column(type: 'uuid', unique: true)]
    private string $id;

    public function getId(): string
    {
        return $this->id;
    }
}

2. Create the batch repository

Create a repository that implements ImportBatchRepositoryInterface and uses ImportBatchRepositoryTrait:

// src/Repository/ImportBatchRepository.php

use Doctrine\Bundle\DoctrineBundle\Repository\ServiceEntityRepository;
use Doctrine\Persistence\ManagerRegistry;
use IQ2i\MessengerImportBundle\Model\ImportBatchRepositoryInterface;
use IQ2i\MessengerImportBundle\Model\ImportBatchRepositoryTrait;

class ImportBatchRepository extends ServiceEntityRepository implements ImportBatchRepositoryInterface
{
    use ImportBatchRepositoryTrait;

    public function __construct(ManagerRegistry $registry)
    {
        parent::__construct($registry, ImportBatch::class);
    }
}

3. Implement BatchAwareMessageInterface on your messages

// src/Message/ImportProductMessage.php

use IQ2i\MessengerImportBundle\Message\BatchAwareMessageInterface;

class ImportProductMessage implements BatchAwareMessageInterface
{
    public function __construct(
        private readonly array $row,
        private readonly ?string $batchId = null,
    ) {}

    public function getRow(): array
    {
        return $this->row;
    }

    public function getBatchId(): ?string
    {
        return $this->batchId;
    }
}

4. Dispatch your messages

Initialize the batch with the total number of messages, then attach the batch ID to each message:

// src/Service/ProductImporter.php

class ProductImporter
{
    public function __construct(
        private readonly ImportBatchRepository $batchRepository,
        private readonly MessageBusInterface $bus,
        private readonly EntityManagerInterface $em,
    ) {}

    public function import(array $rows): void
    {
        $batch = new ImportBatch();
        $batch->initialize(count($rows));

        $this->em->persist($batch);
        $this->em->flush();

        foreach ($rows as $row) {
            $this->bus->dispatch(new ImportProductMessage($row, $batch->getId()));
        }
    }
}

5. Listen to the completion event

// src/EventSubscriber/ImportCompletedSubscriber.php

use IQ2i\MessengerImportBundle\Event\ImportBatchCompletedEvent;
use Symfony\Component\EventDispatcher\EventSubscriberInterface;

class ImportCompletedSubscriber implements EventSubscriberInterface
{
    public static function getSubscribedEvents(): array
    {
        return [
            ImportBatchCompletedEvent::class => 'onImportCompleted',
        ];
    }

    public function onImportCompleted(ImportBatchCompletedEvent $event): void
    {
        // Send a notification, trigger a report, clean up temporary files...
    }
}

API reference

BatchAwareMessageInterface

Implement this interface on any message that belongs to an import batch.

Method Description
getBatchId(): ?string Returns the batch ID, or null if the message is not part of a batch

ImportBatchInterface

Method Description
initialize(int $total): void Sets the total message count and marks the batch as started
getTotal(): int Total number of messages in the batch
getRemaining(): int Number of messages not yet processed
getCreatedAt(): \DateTimeImmutable When the batch was initialized
getCompletedAt(): ?\DateTimeImmutable When the batch completed, null if still in progress
isComplete(): bool Returns true when all messages have been processed
markComplete(): void Sets completedAt (idempotent — safe to call multiple times)

ImportBatchCompletedEvent

Dispatched when the last message in a batch has been handled (or permanently failed).

Method Description
getBatchId(): string ID of the completed batch
getTotal(): int Total number of messages that were processed

Notes

  • Failed messages are counted as processed only when all retry attempts are exhausted (willRetry() === false). A message that will be retried does not decrement the counter.
  • completedAt is set atomically inside ImportBatchRepositoryTrait::decrement() the first time remaining reaches zero. It is safe in concurrent worker environments.
  • The decrement() operation uses a single UPDATE ... WHERE remaining > 0 query to prevent the counter from going below zero under concurrent load.

License

MIT — see LICENSE.

iq2i/messenger-import-bundle 适用场景与选型建议

iq2i/messenger-import-bundle 是一款 基于 PHP 开发的 Composer 扩展包,目前已累计 0 次下载、GitHub Stars 达 0, 最近一次更新时间为 2026 年 03 月 31 日, 在 PHP 生态内属于活跃度较高的组件。

我们在过去多个企业项目中使用过 iq2i/messenger-import-bundle 或与其功能相近的方案,如果你在选型或落地过程中遇到问题,例如 版本兼容、二次改造、私有化封装、与内部系统对接、生产 BUG 排查,欢迎联系我们协助评估。

围绕 iq2i/messenger-import-bundle 我们能提供哪些服务?
定制开发 / 二次开发

基于 iq2i/messenger-import-bundle 在你已有业务上做功能扩展、字段裁剪、UI 适配、与内部账号 / 权限 / 日志系统的深度对接。

BUG 修复 & 性能优化

线上偶发问题、内存泄漏、慢查询、并发异常等排查修复;针对高流量场景做缓存、队列、索引层面的调优。

项目外包 & 长期维护

承接完整的项目从需求 → 设计 → 开发 → 上线 → 长期运维;也可按月提供技术保姆服务。

yvsm@zunyunkeji.com QQ:316430983 微信:yvsm316 西安尊云信息科技 · 专注 PHP / Go / 分布式系统研发

统计信息

  • 总下载量: 0
  • 月度下载量: 0
  • 日度下载量: 0
  • 收藏数: 0
  • 点击次数: 29
  • 依赖项目数: 0
  • 推荐数: 0

GitHub 信息

  • Stars: 0
  • Watchers: 0
  • Forks: 0
  • 开发语言: PHP

其他信息

  • 授权协议: MIT
  • 更新时间: 2026-03-31