Symfony Messenger 异步处理详解
Symfony Messenger 是一个功能强大的消息队列组件,支持异步处理任务,以下是完整的实现指南:

安装和配置
composer require symfony/messenger
config/packages/messenger.yaml:
framework:
messenger:
# 失败消息处理
failure_transport: failed_default
transports:
# 异步传输(使用 Doctrine 作为存储)
async:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
use_notify: true
check_delayed_interval: 60000
retry_strategy:
max_retries: 3
delay: 1000
multiplier: 2
max_delay: 0
# 同步传输(用于调试)
sync: 'sync://'
# 失败消息存储
failed_default: 'doctrine://default?queue_name=failed'
routing:
# 将特定消息路由到异步传输
'App\Message\SendEmailMessage': async
'App\Message\ProcessPaymentMessage': async
# 总线下使用中间件
buses:
messenger.bus.default:
middleware:
- doctrine_transaction
- validation
- logger
创建消息类
// src/Message/SendEmailMessage.php
namespace App\Message;
class SendEmailMessage
{
public function __construct(
private string $email,
private string $subject,
private string $content,
private array $options = []
) {}
public function getEmail(): string
{
return $this->email;
}
public function getSubject(): string
{
return $this->subject;
}
public function getContent(): string
{
return $this->content;
}
public function getOptions(): array
{
return $this->options;
}
}
创建消息处理器
// src/MessageHandler/SendEmailHandler.php
namespace App\MessageHandler;
use App\Message\SendEmailMessage;
use Symfony\Component\Mailer\MailerInterface;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;
use Symfony\Component\Messenger\Exception\UnrecoverableMessageHandlingException;
use Psr\Log\LoggerInterface;
#[AsMessageHandler]
class SendEmailHandler
{
public function __construct(
private MailerInterface $mailer,
private LoggerInterface $logger
) {}
public function __invoke(SendEmailMessage $message): void
{
try {
$this->logger->info('Sending email', [
'email' => $message->getEmail(),
'subject' => $message->getSubject()
]);
// 实际的邮件发送逻辑
$email = (new Email())
->from('sender@example.com')
->to($message->getEmail())
->subject($message->getSubject())
->text($message->getContent());
$this->mailer->send($email);
} catch (\Exception $e) {
$this->logger->error('Failed to send email', [
'error' => $e->getMessage()
]);
// 抛出不可恢复异常,防止重试
throw new UnrecoverableMessageHandlingException(
'Failed to send email: ' . $e->getMessage()
);
}
}
}
发送消息
// src/Controller/EmailController.php
namespace App\Controller;
use App\Message\SendEmailMessage;
use Symfony\Bundle\FrameworkBundle\Controller\AbstractController;
use Symfony\Component\HttpFoundation\Response;
use Symfony\Component\Messenger\MessageBusInterface;
use Symfony\Component\Routing\Annotation\Route;
class EmailController extends AbstractController
{
#[Route('/send-email', name: 'send_email')]
public function sendEmail(MessageBusInterface $bus): Response
{
// 立即发送消息到队列
$bus->dispatch(new SendEmailMessage(
email: 'user@example.com',
subject: 'Welcome!',
content: 'Thank you for registering.'
));
return new Response('Email queued for sending!');
}
}
配置传输 DSN
.env 或 .env.local:
# 使用 Doctrine MESSENGER_TRANSPORT_DSN=doctrine://default # 使用 Redis MESSENGER_TRANSPORT_DSN=redis://localhost:6379/messages # 使用 AMQP (RabbitMQ) MESSENGER_TRANSPORT_DSN=amqp://guest:guest@localhost:5672/%2f/messages # 使用 Amazon SQS MESSENGER_TRANSPORT_DSN=sqs://sqs.us-east-1.amazonaws.com/your-queue # 使用文件系统(开发环境) MESSENGER_TRANSPORT_DSN=file:///var/tmp/messages
启动消费者
# 基本用法 php bin/console messenger:consume async -vv # 指定最大消息处理数量 php bin/console messenger:consume async --limit=10 # 指定运行时间(秒) php bin/console messenger:consume async --time-limit=60 # 指定内存限制(MB) php bin/console messenger:consume async --memory-limit=128 # 后台运行 nohup php bin/console messenger:consume async &> messenger.log & # 使用 Supervisor 管理(推荐)
Supervisor 配置
/etc/supervisor/conf.d/messenger-worker.conf:
[program:messenger-worker] process_name=%(program_name)s_%(process_num)02d command=php /path/to/your/project/bin/console messenger:consume async --time-limit=3600 --memory-limit=256M numprocs=4 autostart=true autorestart=true user=www-data redirect_stderr=true stdout_logfile=/path/to/your/project/var/log/messenger.log stopwaitsecs=3600
高级特性
1 消息中间件
// config/packages/messenger.yaml
framework:
messenger:
buses:
messenger.bus.default:
middleware:
# 自定义中间件
- App\Middleware\AuditMiddleware
- doctrine_transaction
- validation
// src/Middleware/AuditMiddleware.php
namespace App\Middleware;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Middleware\MiddlewareInterface;
use Symfony\Component\Messenger\Middleware\StackInterface;
use Psr\Log\LoggerInterface;
class AuditMiddleware implements MiddlewareInterface
{
public function __construct(
private LoggerInterface $logger
) {}
public function handle(Envelope $envelope, StackInterface $stack): Envelope
{
$message = $envelope->getMessage();
$this->logger->info('Processing message', [
'class' => get_class($message),
'time' => date('Y-m-d H:i:s')
]);
$start = microtime(true);
try {
$envelope = $stack->next()->handle($envelope, $stack);
$this->logger->info('Message processed successfully', [
'duration' => microtime(true) - $start
]);
return $envelope;
} catch (\Exception $e) {
$this->logger->error('Message processing failed', [
'error' => $e->getMessage(),
'duration' => microtime(true) - $start
]);
throw $e;
}
}
}
2 延迟消息
// 延迟10秒处理
use Symfony\Component\Messenger\Stamp\DelayStamp;
$bus->dispatch(
new SendEmailMessage('user@example.com', 'Subject', 'Content'),
[new DelayStamp(10000)] // 10秒延迟
);
3 消息优先级
framework:
messenger:
transports:
high_priority:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
queue_name: high_priority
low_priority:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
queue_name: low_priority
routing:
'App\Message\UrgentMessage': high_priority
'App\Message\NormalMessage': [async, low_priority]
错误处理和重试
// 处理失败消息 php bin/console messenger:failed:show # 查看失败消息 php bin/console messenger:failed:retry # 重试失败消息 php bin/console messenger:failed:remove # 删除失败消息
监控和调试
framework:
messenger:
transports:
async:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
auto_setup: true
use_notify: true
check_delayed_interval: 60000
查看消息状态:
# 监控消费者状态 php bin/console debug:messenger # 查看传输统计 php bin/console doctrine:query:sql 'SELECT * FROM messenger_messages'
最佳实践
- 消息设计:保持消息小巧,只传递必要数据
- 幂等性:确保消息处理具有幂等性
- 异常处理:合理使用
UnrecoverableMessageHandlingException - 监控:实现完善的日志和监控机制
- 资源管理:使用 Supervisor 管理消费者进程
这个框架使得异步处理变得简单而强大,适合处理邮件发送、数据导入、图片处理等耗时的后台任务。