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
相关产品推荐
相关产品推荐

