ProcessDeadLetterQueue.php 2.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293
  1. <?php
  2. namespace App\Console\Commands;
  3. use Illuminate\Console\Command;
  4. use PhpAmqpLib\Connection\AMQPStreamConnection;
  5. use PhpAmqpLib\Message\AMQPMessage;
  6. class ProcessDeadLetterQueue extends Command
  7. {
  8. /**
  9. * The name and signature of the console command.
  10. * 查看死信队列消息
  11. * php artisan rabbitmq:process-dlq orders_dlq
  12. *
  13. * 重新入队死信消息
  14. * php artisan rabbitmq:process-dlq orders_dlq --requeue
  15. *
  16. * 删除死信消息
  17. * php artisan rabbitmq:process-dlq orders_dlq --delete
  18. *
  19. * @var string
  20. */
  21. protected $signature = 'rabbitmq:process-dlq {dlq_name} {--requeue} {--delete}';
  22. protected $description = '处理死信队列中的消息';
  23. public function handle()
  24. {
  25. $dlqName = $this->argument('dlq_name');
  26. $requeue = $this->option('requeue');
  27. $delete = $this->option('delete');
  28. $config = config('queue.connections.rabbitmq');
  29. $connection = new AMQPStreamConnection(
  30. $config['host'],
  31. $config['port'],
  32. $config['user'],
  33. $config['password'],
  34. $config['virtual_host']
  35. );
  36. $channel = $connection->channel();
  37. $this->info("开始处理死信队列: {$dlqName}");
  38. $messageCount = 0;
  39. while (true) {
  40. $msg = $channel->basic_get($dlqName, false);
  41. if (! $msg) {
  42. break; // 队列为空
  43. }
  44. $messageCount++;
  45. $data = json_decode($msg->body, true);
  46. $this->info("处理第 {$messageCount} 条死信消息");
  47. $this->line('原始队列: '.($data['queue'] ?? 'unknown'));
  48. $this->line('失败原因: '.($data['failure_reason'] ?? 'unknown'));
  49. $this->line('失败时间: '.($data['failed_at'] ?? 'unknown'));
  50. if ($requeue) {
  51. // 重新入队到原始队列
  52. $originalQueue = $data['queue'] ?? null;
  53. if ($originalQueue && isset($data['original_message'])) {
  54. $requeueMsg = new AMQPMessage(
  55. json_encode($data['original_message']),
  56. ['delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT]
  57. );
  58. $channel->basic_publish($requeueMsg, '', $originalQueue);
  59. $this->info("消息已重新入队到: {$originalQueue}");
  60. }
  61. }
  62. if ($delete || $requeue) {
  63. $msg->ack();
  64. } else {
  65. // 只是查看,不删除
  66. $msg->nack(false, true);
  67. }
  68. }
  69. $this->info("死信队列处理完成,共处理 {$messageCount} 条消息");
  70. $channel->close();
  71. $connection->close();
  72. return 0;
  73. }
  74. }