如何在Symfony Messenger总线中批量处理消息并合并报告邮件
基于Symfony Messenger实现合并报告邮件的解决方案
要实现批量生成报告后发送合并邮件的需求,核心是跟踪所有报告任务的完成状态(无论是否生成成功),待全部任务结束后触发合并邮件发送。以下是具体的实现方案:
一、核心思路:任务协调与状态跟踪
采用「批次协调器+任务状态跟踪」的模式,通过以下步骤实现:
- 生成唯一批次ID,记录该批次的总任务数
- 为每个报告生成单独的异步任务,关联批次ID
- 每个任务完成后(成功或失败),更新批次的完成计数
- 当完成计数等于总任务数时,收集所有成功生成的报告,发送合并邮件
二、具体实现步骤
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); } } }
三、关键细节与边缘情况处理
- 任务失败/无输出处理:无论报告是否生成成功,只要任务执行完成(包括失败),都会更新完成计数,确保批次能正常触发合并邮件,仅附加成功生成的报告。
- 任务超时处理:可添加定时任务(如Cron),检查超过一定时间(如2小时)仍未完成的批次,强制触发邮件并标记未完成任务为失败。
- 幂等性保障:启用Symfony Messenger的消息去重功能,避免重复处理任务完成通知。在
config/packages/messenger.yaml中配置:framework: messenger: transports: async: dsn: '%env(MESSENGER_TRANSPORT_DSN)%' options: deduplicate: true - 临时文件管理:报告文件建议存储在
var/tmp/reports目录下,确保Messenger worker有权限访问,且发送邮件后及时清理。 - 存储选择:示例使用Redis作为状态存储(性能更高),若没有Redis,可改用数据库表存储批次状态(字段包括
batch_id、total_tasks、completed_tasks、reports(JSON格式))。
内容的提问来源于stack exchange,提问作者Jason Olson
相关产品推荐
相关产品推荐

