BaseRabbitMQJob.php 3.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124
  1. <?php
  2. namespace App\Jobs;
  3. use App\Exceptions\TaskFailException;
  4. use Illuminate\Bus\Queueable;
  5. use Illuminate\Contracts\Queue\ShouldQueue;
  6. use Illuminate\Foundation\Bus\Dispatchable;
  7. use Illuminate\Queue\InteractsWithQueue;
  8. use Illuminate\Queue\SerializesModels;
  9. use Illuminate\Support\Facades\Log;
  10. abstract class BaseRabbitMQJob implements ShouldQueue
  11. {
  12. use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
  13. protected $queueName;
  14. protected $messageData;
  15. protected $currentRetryCount = 0;
  16. protected $tries = 0;
  17. protected $timeout = 0;
  18. protected $messageId = null;
  19. protected $stop = false;
  20. public function __construct(string $queueName, string $messageId, array $messageData, int $retryCount = 0)
  21. {
  22. $this->queueName = $queueName;
  23. $this->messageData = $messageData;
  24. $this->currentRetryCount = $retryCount;
  25. $this->messageId = $messageId;
  26. // 从配置读取重试次数和超时时间
  27. $queueConfig = config("mint.rabbitmq.queues.{$queueName}");
  28. $this->tries = $queueConfig['retry_times'] ?? 3;
  29. $this->timeout = $queueConfig['timeout'] ?? 300;
  30. }
  31. public function handle()
  32. {
  33. try {
  34. Log::info('开始处理队列消息', [
  35. 'queue' => $this->queueName,
  36. 'message_id' => $this->messageId ?? 'unknown',
  37. 'retry_count' => $this->currentRetryCount,
  38. ]);
  39. // 调用子类的具体业务逻辑
  40. $result = $this->processMessage($this->messageData);
  41. Log::info('队列消息处理完成', [
  42. 'queue' => $this->queueName,
  43. 'message_id' => $this->messageId ?? 'unknown',
  44. 'result' => $result,
  45. ]);
  46. return $result;
  47. } catch (TaskFailException $e) {
  48. $this->handleFinalFailure($this->messageData, $e);
  49. throw $e;
  50. } catch (\Exception $e) {
  51. Log::error('队列消息处理失败', [
  52. 'queue' => $this->queueName,
  53. 'message_id' => $this->messageId ?? 'unknown',
  54. 'error' => $e->getMessage(),
  55. 'retry_count' => $this->currentRetryCount,
  56. 'max_retries' => $this->tries,
  57. ]);
  58. // 如果达到最大重试次数,处理失败逻辑
  59. if ($this->currentRetryCount >= $this->tries - 1) {
  60. $this->handleFinalFailure($this->messageData, $e);
  61. }
  62. throw $e; // 重新抛出异常以触发重试
  63. }
  64. }
  65. public function failed(\Exception $exception)
  66. {
  67. Log::error('队列消息最终失败', [
  68. 'queue' => $this->queueName,
  69. 'message_id' => $this->messageId ?? 'unknown',
  70. 'error' => $exception->getMessage(),
  71. 'retry_count' => $this->currentRetryCount,
  72. ]);
  73. // 发送到死信队列的逻辑将在 Worker 中处理
  74. }
  75. // 子类需要实现的具体业务逻辑
  76. abstract protected function processMessage(array $messageData);
  77. // 子类可以重写的失败处理逻辑
  78. protected function handleFinalFailure(array $messageData, \Exception $exception)
  79. {
  80. // 默认实现:记录日志
  81. Log::error('消息处理最终失败,准备发送到死信队列', [
  82. 'queue' => $this->queueName,
  83. 'message_id' => $this->messageId ?? 'unknown',
  84. 'error' => $exception->getMessage(),
  85. ]);
  86. }
  87. public function getQueueName(): string
  88. {
  89. return $this->queueName;
  90. }
  91. public function getCurrentRetryCount(): int
  92. {
  93. return $this->currentRetryCount;
  94. }
  95. public function stop()
  96. {
  97. $this->stop = true;
  98. }
  99. }