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

Laravel中创建可中断续跑的分批更新命令及Job实现方案

Laravel 分批更新数据并支持中断续跑的实现方案

现有代码的问题

你的当前实现存在几个关键问题:

  • 批次范围计算逻辑错误:当总数据量大于100时,intval($woocount/100)得到的是批次数而非批次结束ID,导致后续ID范围查询完全错误。
  • 依赖连续ID:如果数据库中存在ID不连续的情况(比如删除过数据),会出现漏处理或批次数据为空的情况。
  • 无状态记录:没有保存批次处理状态,中断后重启会从头执行,无法续跑。
  • 直接传递集合到Job:可能引发序列化问题,且大集合会占用过多内存。

正确实现方案

我们通过批次状态记录+分页查询来解决问题,核心思路是先生成所有待处理的批次记录,Job处理时基于批次记录执行更新,同时维护批次状态,确保中断后仅处理未完成的批次。

1. 创建批次记录模型与迁移

首先创建用于记录批次状态的模型和数据库表:

生成迁移文件:

php artisan make:migration create_sync_batches_table

修改迁移文件内容:

use Illuminate\Database\Migrations\Migration;
use Illuminate\Database\Schema\Blueprint;
use Illuminate\Support\Facades\Schema;

return new class extends Migration
{
    public function up()
    {
        Schema::create('sync_batches', function (Blueprint $table) {
            $table->id();
            $table->integer('batch_number')->unique(); // 批次号
            $table->unsignedBigInteger('start_id'); // 批次起始ID
            $table->unsignedBigInteger('end_id'); // 批次结束ID
            $table->enum('status', ['pending', 'processing', 'completed', 'failed'])->default('pending'); // 批次状态
            $table->integer('total_items'); // 批次总数据量
            $table->integer('processed_items')->default(0); // 已处理数量
            $table->timestamps();
        });
    }

    public function down()
    {
        Schema::dropIfExists('sync_batches');
    }
};

运行迁移:

php artisan migrate

创建模型:

php artisan make:model SyncBatch

2. 编写Artisan命令

创建一个生成批次并调度Job的命令:

php artisan make:command SyncProductsStock

修改命令文件内容:

namespace App\Console\Commands;

use App\Models\SyncBatch;
use App\Models\WoocommerceProduct;
use App\Jobs\CheckStockJob;
use Illuminate\Console\Command;

class SyncProductsStock extends Command
{
    protected $signature = 'products:sync-stock';
    protected $description = '分批同步商品库存到API,支持中断续跑';

    private const BATCH_SIZE = 50; // 每批次50条

    public function handle()
    {
        $this->info('开始生成待处理批次...');

        // 1. 获取所有需要同步的商品ID范围
        $minId = WoocommerceProduct::where('sync_status', 'IN_SYNC')->min('id');
        $maxId = WoocommerceProduct::where('sync_status', 'IN_SYNC')->max('id');

        if (!$minId || !$maxId) {
            $this->info('没有需要同步的商品');
            return;
        }

        // 2. 生成所有未创建的批次记录
        $currentBatchNumber = 1;
        $currentStartId = $minId;

        while ($currentStartId <= $maxId) {
            $currentEndId = $currentStartId + self::BATCH_SIZE - 1;
            // 避免最后一批次超出最大ID
            if ($currentEndId > $maxId) {
                $currentEndId = $maxId;
            }

            // 检查该批次是否已存在
            $existingBatch = SyncBatch::where('start_id', $currentStartId)
                ->where('end_id', $currentEndId)
                ->first();

            if (!$existingBatch) {
                // 计算该批次的实际数据量
                $totalItems = WoocommerceProduct::where('sync_status', 'IN_SYNC')
                    ->whereBetween('id', [$currentStartId, $currentEndId])
                    ->count();

                SyncBatch::create([
                    'batch_number' => $currentBatchNumber,
                    'start_id' => $currentStartId,
                    'end_id' => $currentEndId,
                    'total_items' => $totalItems,
                ]);
            }

            $currentStartId = $currentEndId + 1;
            $currentBatchNumber++;
        }

        $this->info('批次生成完成,开始调度处理任务...');

        // 3. 调度所有未完成的批次Job
        $pendingBatches = SyncBatch::whereIn('status', ['pending', 'failed'])->get();

        foreach ($pendingBatches as $batch) {
            CheckStockJob::dispatch($batch->id);
            $this->info('已调度批次: ' . $batch->batch_number);
        }

        $this->info('所有待处理批次已调度完成');
    }
}

3. 编写处理批次的Job

修改CheckStockJob,基于批次记录处理数据:

php artisan make:job CheckStockJob

修改Job文件内容:

namespace App\Jobs;

use App\Models\SyncBatch;
use App\Models\WoocommerceProduct;
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\Http;

class CheckStockJob implements ShouldQueue
{
    use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;

    private $batchId;

    public function __construct(int $batchId)
    {
        $this->batchId = $batchId;
    }

    public function handle()
    {
        $batch = SyncBatch::findOrFail($this->batchId);

        // 标记批次为处理中
        $batch->update(['status' => 'processing']);

        try {
            // 查询当前批次的商品数据,cursor()减少内存占用
            $products = WoocommerceProduct::where('sync_status', 'IN_SYNC')
                ->whereBetween('id', [$batch->start_id, $batch->end_id])
                ->cursor();

            $processedCount = 0;

            foreach ($products as $product) {
                // 调用API更新数据,替换为你的实际API逻辑
                $response = Http::put('https://your-api-url.com/products/' . $product->external_id, [
                    'stock_quantity' => $product->stock_quantity,
                    // 其他需要更新的字段
                ]);

                if ($response->successful()) {
                    // 可按需更新商品同步状态
                    // $product->update(['sync_status' => 'SYNCED']);
                    $processedCount++;
                } else {
                    throw new \Exception('API更新失败: ' . $response->body());
                }
            }

            // 标记批次为已完成
            $batch->update([
                'status' => 'completed',
                'processed_items' => $processedCount,
            ]);

        } catch (\Exception $e) {
            // 标记批次为失败,可添加error_message字段记录错误
            $batch->update([
                'status' => 'failed',
                // 'error_message' => $e->getMessage(),
            ]);

            // 抛出异常触发队列重试(可在Job中配置重试次数)
            throw $e;
        }
    }
}

4. 关键优势说明

  • 中断续跑支持:通过SyncBatch记录批次状态,重启命令后只会调度未完成(pending/failed)的批次,不会重复处理已完成的批次。
  • 内存优化:使用cursor()而非get()查询数据,避免一次性加载大量数据到内存。
  • 可靠的批次划分:基于ID范围+实际数据量统计,即使ID不连续也能确保每个批次的准确性。
  • 状态可视化:可通过sync_batches表查看每个批次的处理进度、状态,方便排查问题。

5. 使用方式

运行命令开始同步:

php artisan products:sync-stock

如果需要定时执行,可在app/Console/Kernel.php中添加Cron任务:

protected function schedule(Schedule $schedule)
{
    $schedule->command('products:sync-stock')->daily(); // 每天执行一次
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 12:50:22