Laravel 11调度消息命令异常:重复使用旧值问题求助
Laravel 11调度消息命令重复分发问题
在Laravel 11中实现了发送调度消息的控制台命令,逻辑是从数据库获取待发送的调度消息,分发SendScheduledMessageJob处理。但目前出现异常:第一条消息处理正常,后续消息重复使用上一次循环的数值,导致消息分发错误。
命令代码
<?php namespace App\Console\Commands; use App\Jobs\SendScheduledMessageJob; use App\Models\ScheduledMessage; use Illuminate\Console\Command; class SendScheduledMessages extends Command { /** * The name and signature of the console command. * * @var string */ protected $signature = 'app:send-scheduled-messages'; /** * The console command description. * * @var string */ protected $description = 'Send scheduled messages'; /** * Execute the console command. */ public function handle() { $messages = ScheduledMessage::where('scheduled_at', '<=', now()) ->where('is_sent', false) ->get(); foreach ($messages as $message) { SendScheduledMessageJob::dispatch($message); } } }
调度任务路由
<?php use App\Models\ScheduledMessage; use Illuminate\Foundation\Inspiring; use Illuminate\Support\Facades\Artisan; use Illuminate\Support\Facades\Log; use Illuminate\Support\Facades\Schedule; use Illuminate\Support\Facades\Storage; Artisan::command('inspire', function () { $this->comment(Inspiring::quote()); })->purpose('Display an inspiring quote')->hourly(); Schedule::command('app:send-scheduled-messages')->everySecond()->description('Send scheduled messages');
SendScheduledMessageJob代码
<?php namespace App\Jobs; use App\Models\ScheduledMessage; use Illuminate\Bus\Queueable; use Illuminate\Contracts\Queue\ShouldQueue; use Illuminate\Foundation\Bus\Dispatchable; use Illuminate\Queue\InteractsWithQueue; use Illuminate\Queue\SerializesModels; use Illuminate\Support\Facades\DB; use Illuminate\Support\Facades\Log; use Illuminate\Support\Facades\Storage; class SendScheduledMessageJob implements ShouldQueue { use Dispatchable, InteractsWithQueue, Queueable, SerializesModels; protected $message; /** * Create a new job instance. */ public function __construct(ScheduledMessage $message) { $this->message = $message; Log::info($this->message->scheduled_at); } /** * Execute the job. */ public function handle(): void { DB::transaction(function () { try { // Retrieve the messageable entity (the recipient) $messageable = $this->message->messageable; // Create a new message for the recipient with the content and sender ID $newMessage = $messageable->messages()->create([ 'content' => $this->message->content, 'sender_id' => $this->message->sender_id, ]); // Check if the original message has attachments if ($this->message->attachments->isNotEmpty()) { // Iterate through each attachment foreach ($this->message->attachments as $attachment) { // Generate a new path for the attachment, replacing the old message ID with the new one $path = str_replace($this->message->id, $newMessage->id, $attachment->path); // Determine the directory to remove after moving the file, up to "/attachments" $path_to_remove = substr($attachment->path, 0, strpos($attachment->path, "/attachments")); // Move the attachment to the new path Storage::disk('public')->move($attachment->path, $path); // Create a record for the attachment under the new message quietly (without raising events) $newMessage->attachments()->createQuietly([ 'path' => $path, 'name' => $attachment->name, ]); // Delete the old directory where the attachment was stored Storage::disk('public')->deleteDirectory($path_to_remove); } } // Mark the original message as sent $this->message->update(['is_sent' => true]); } catch (\Exception $e) { // Log any exceptions Log::error('Failed to send message ID: ' . $this->message->id . ', Error: ' . $e->getMessage()); } }); } }
数据库中的调度消息时间
- 2024-07-03 12:14:00
- 2024-07-03 12:15:00
- 2024-07-03 12:16:00
- 2024-07-03 12:17:00
生成的日志
[2024-07-03 14:08:00] local.INFO: 2024-07-03 12:14:00 [2024-07-03 14:08:00] local.INFO: 2024-07-03 12:14:00 [2024-07-03 14:08:00] local.INFO: 2024-07-03 12:15:00 [2024-07-03 14:08:00] local.INFO: 2024-07-03 12:16:00
解决思路
避免重复获取未处理消息
当前调度每秒执行一次,可能Job还未完成is_sent状态更新时,下一次调度又获取到同一条消息。可以新增is_processing字段,先标记消息为处理中再获取:// 先标记待处理消息 ScheduledMessage::where('scheduled_at', '<=', now()) ->where('is_sent', false) ->where('is_processing', false) ->update(['is_processing' => true]); // 获取已标记的消息 $messages = ScheduledMessage::where('is_processing', true)->get();处理完成后同时更新
is_sent和is_processing为true和false。改用模型ID传递而非实例
避免模型序列化带来的实例复用问题,分发Job时只传ID,在Job内部重新查询最新数据:// 命令中 foreach ($messages as $message) { SendScheduledMessageJob::dispatch($message->id); } // Job中 public function __construct(int $messageId) { $this->messageId = $messageId; } public function handle(): void { $message = ScheduledMessage::findOrFail($this->messageId); // 后续逻辑使用$message }调整调度频率
每秒执行过于频繁,容易引发并发冲突,可根据业务需求调整为每30秒/1分钟执行一次。使用悲观锁防止并发
在Job处理时给消息加锁,确保同一时间只有一个进程处理该消息:$message = ScheduledMessage::where('id', $this->message->id)->lockForUpdate()->first(); // 执行处理逻辑
内容的提问来源于stack exchange,提问作者Najeeb Anwari
相关产品推荐
相关产品推荐

