You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

解决思路

  1. 避免重复获取未处理消息
    当前调度每秒执行一次,可能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。

  2. 改用模型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
    }
    
  3. 调整调度频率
    每秒执行过于频繁,容易引发并发冲突,可根据业务需求调整为每30秒/1分钟执行一次。

  4. 使用悲观锁防止并发
    在Job处理时给消息加锁,确保同一时间只有一个进程处理该消息:

    $message = ScheduledMessage::where('id', $this->message->id)->lockForUpdate()->first();
    // 执行处理逻辑
    

内容的提问来源于stack exchange,提问作者Najeeb Anwari

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 15:54:54