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

如何在Symfony Messenger总线中批量处理消息并合并报告邮件

基于Symfony Messenger实现合并报告邮件的解决方案

要实现批量生成报告后发送合并邮件的需求,核心是跟踪所有报告任务的完成状态(无论是否生成成功),待全部任务结束后触发合并邮件发送。以下是具体的实现方案:

一、核心思路:任务协调与状态跟踪

采用「批次协调器+任务状态跟踪」的模式,通过以下步骤实现:

  1. 生成唯一批次ID,记录该批次的总任务数
  2. 为每个报告生成单独的异步任务,关联批次ID
  3. 每个任务完成后(成功或失败),更新批次的完成计数
  4. 当完成计数等于总任务数时,收集所有成功生成的报告,发送合并邮件

二、具体实现步骤

1. 定义所需消息类

创建三个核心消息类,用于任务分发、状态通知和邮件触发:

// src/Message/GenerateSingleReport.php
namespace App\Message;

class GenerateSingleReport
{
    public function __construct(
        private string $batchId, // 批次ID
        private array $reportParams // 报告生成参数
    ) {}

    // Getter方法
    public function getBatchId(): string { return $this->batchId; }
    public function getReportParams(): array { return $this->reportParams; }
}

// src/Message/ReportTaskCompleted.php
namespace App\Message;

class ReportTaskCompleted
{
    public function __construct(
        private string $batchId,
        private bool $isSuccess,
        private ?string $reportPath = null // 成功时的报告文件路径
    ) {}

    // Getter方法
    public function getBatchId(): string { return $this->batchId; }
    public function isSuccess(): bool { return $this->isSuccess; }
    public function getReportPath(): ?string { return $this->reportPath; }
}

// src/Message/SendMergedReportEmail.php
namespace App\Message;

class SendMergedReportEmail
{
    public function __construct(
        private string $batchId,
        private string $recipient,
        private array $reportPaths // 所有成功生成的报告路径
    ) {}

    // Getter方法
    public function getRecipient(): string { return $this->recipient; }
    public function getReportPaths(): array { return $this->reportPaths; }
}

2. 初始化批次并分发任务

查询数据库获取需生成的报告列表后,初始化批次信息并分发每个报告的生成任务:

// 假设在某个服务或控制器中
use App\Message\GenerateSingleReport;
use Symfony\Component\Messenger\MessageBusInterface;
use Symfony\Component\Uid\Uuid;
use Symfony\Component\Cache\Adapter\RedisAdapter;

public function dispatchReportTasks(array $reportItems, string $recipient, MessageBusInterface $bus, RedisAdapter $redis): void
{
    $batchId = Uuid::v4()->toString(); // 生成唯一批次ID
    $totalTasks = count($reportItems);

    // 初始化Redis中的批次状态
    $redis->set("batch:$batchId:total", $totalTasks);
    $redis->set("batch:$batchId:completed", 0);
    $redis->del("batch:$batchId:reports"); // 清空旧数据

    // 保存收件人信息
    $redis->set("batch:$batchId:recipient", $recipient);

    // 分发每个报告生成任务
    foreach ($reportItems as $params) {
        $bus->dispatch(new GenerateSingleReport($batchId, $params));
    }
}

3. 处理单个报告生成任务

报告生成完成后,无论成功或失败,都分发任务完成的通知消息:

// src/MessageHandler/GenerateSingleReportHandler.php
namespace App\MessageHandler;

use App\Message\GenerateSingleReport;
use App\Message\ReportTaskCompleted;
use Symfony\Component\Messenger\MessageBusInterface;
use Psr\Log\LoggerInterface;

class GenerateSingleReportHandler
{
    public function __construct(
        private MessageBusInterface $bus,
        private YourReportGeneratorService $reportGenerator,
        private LoggerInterface $logger
    ) {}

    public function __invoke(GenerateSingleReport $message): void
    {
        $reportPath = null;
        $isSuccess = false;

        try {
            // 调用现有报告生成逻辑
            $reportPath = $this->reportGenerator->generate($message->getReportParams());
            $isSuccess = true;
        } catch (\Exception $e) {
            // 记录错误,但仍标记任务完成
            $this->logger->error('报告生成失败', [
                'batchId' => $message->getBatchId(),
                'params' => $message->getReportParams(),
                'error' => $e->getMessage()
            ]);
        }

        // 分发任务完成通知
        $this->bus->dispatch(new ReportTaskCompleted(
            $message->getBatchId(),
            $isSuccess,
            $reportPath
        ));
    }
}

4. 跟踪批次完成状态并触发合并邮件

处理任务完成通知,更新批次状态,当所有任务完成时触发邮件发送:

// src/MessageHandler/ReportTaskCompletedHandler.php
namespace App\MessageHandler;

use App\Message\ReportTaskCompleted;
use App\Message\SendMergedReportEmail;
use Symfony\Component\Messenger\MessageBusInterface;
use Symfony\Component\Cache\Adapter\RedisAdapter;

class ReportTaskCompletedHandler
{
    public function __construct(
        private MessageBusInterface $bus,
        private RedisAdapter $redis
    ) {}

    public function __invoke(ReportTaskCompleted $message): void
    {
        $batchId = $message->getBatchId();
        $keys = [
            'total' => "batch:$batchId:total",
            'completed' => "batch:$batchId:completed",
            'reports' => "batch:$batchId:reports",
            'recipient' => "batch:$batchId:recipient"
        ];

        // 增加完成计数
        $completedCount = $this->redis->increment($keys['completed']);
        // 获取总任务数
        $totalCount = (int)$this->redis->get($keys['total']);

        // 如果任务成功,保存报告路径
        if ($message->isSuccess() && $message->getReportPath()) {
            $this->redis->lPush($keys['reports'], $message->getReportPath());
        }

        // 检查是否所有任务已完成
        if ($completedCount === $totalCount) {
            // 获取所有成功生成的报告路径
            $reportPaths = $this->redis->lRange($keys['reports'], 0, -1);
            // 获取收件人邮箱
            $recipient = $this->redis->get($keys['recipient']);

            // 触发合并邮件发送任务
            $this->bus->dispatch(new SendMergedReportEmail(
                $batchId,
                $recipient,
                $reportPaths
            ));

            // 设置Redis键的过期时间(24小时后自动清理)
            foreach ($keys as $key) {
                $this->redis->expire($key, 86400);
            }
        }
    }
}

5. 发送合并邮件并清理临时文件

处理邮件发送任务,附加所有报告文件并清理临时文件:

// src/MessageHandler/SendMergedReportEmailHandler.php
namespace App\MessageHandler;

use App\Message\SendMergedReportEmail;
use Symfony\Component\Mailer\MailerInterface;
use Symfony\Component\Mime\Email;
use Symfony\Component\Filesystem\Filesystem;

class SendMergedReportEmailHandler
{
    public function __construct(
        private MailerInterface $mailer,
        private Filesystem $filesystem
    ) {}

    public function __invoke(SendMergedReportEmail $message): void
    {
        $email = (new Email())
            ->from('reports@yourdomain.com')
            ->to($message->getRecipient())
            ->subject('您的合并报告')
            ->text('请查收附件中的所有生成报告。');

        // 附加所有报告文件
        foreach ($message->getReportPaths() as $path) {
            if ($this->filesystem->exists($path)) {
                $email->attachFromPath($path);
            }
        }

        // 发送邮件
        $this->mailer->send($email);

        // 清理临时报告文件
        foreach ($message->getReportPaths() as $path) {
            $this->filesystem->remove($path);
        }
    }
}

三、关键细节与边缘情况处理

  1. 任务失败/无输出处理:无论报告是否生成成功,只要任务执行完成(包括失败),都会更新完成计数,确保批次能正常触发合并邮件,仅附加成功生成的报告。
  2. 任务超时处理:可添加定时任务(如Cron),检查超过一定时间(如2小时)仍未完成的批次,强制触发邮件并标记未完成任务为失败。
  3. 幂等性保障:启用Symfony Messenger的消息去重功能,避免重复处理任务完成通知。在config/packages/messenger.yaml中配置:
    framework:
        messenger:
            transports:
                async:
                    dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
                    options:
                        deduplicate: true
    
  4. 临时文件管理:报告文件建议存储在var/tmp/reports目录下,确保Messenger worker有权限访问,且发送邮件后及时清理。
  5. 存储选择:示例使用Redis作为状态存储(性能更高),若没有Redis,可改用数据库表存储批次状态(字段包括batch_id、total_tasks、completed_tasks、reports(JSON格式))。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 18:27:02