PHP项目Symfony messenger与异步

wen PHP项目 4

Symfony Messenger 异步处理详解

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

PHP项目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'

最佳实践

  1. 消息设计:保持消息小巧,只传递必要数据
  2. 幂等性:确保消息处理具有幂等性
  3. 异常处理:合理使用 UnrecoverableMessageHandlingException
  4. 监控:实现完善的日志和监控机制
  5. 资源管理:使用 Supervisor 管理消费者进程

这个框架使得异步处理变得简单而强大,适合处理邮件发送、数据导入、图片处理等耗时的后台任务。

抱歉,评论功能暂时关闭!