|
@@ -4,12 +4,12 @@ namespace App\Console\Commands;
|
|
|
|
|
|
|
|
use Illuminate\Console\Command;
|
|
use Illuminate\Console\Command;
|
|
|
use Illuminate\Support\Facades\DB;
|
|
use Illuminate\Support\Facades\DB;
|
|
|
-use App\Services\DeepSeek\DeepSeekService;
|
|
|
|
|
-use App\Facade\Site;
|
|
|
|
|
|
|
+use Symfony\Component\Process\Process;
|
|
|
|
|
|
|
|
/**
|
|
/**
|
|
|
- * 批量生成分集定时任务
|
|
|
|
|
- * 每次执行处理一个待生成的剧集
|
|
|
|
|
|
|
+ * 批量生成分集定时任务(调度器)
|
|
|
|
|
+ * 每次执行:找出有 pending 任务的 anime_id,为最多 N 个(默认5)空闲 anime 各拉起一个后台子进程,
|
|
|
|
|
+ * 由子进程独占处理对应 anime 的待生成剧集。
|
|
|
*/
|
|
*/
|
|
|
class ProcessBatchEpisodeGenerationCommand extends Command
|
|
class ProcessBatchEpisodeGenerationCommand extends Command
|
|
|
{
|
|
{
|
|
@@ -25,20 +25,7 @@ class ProcessBatchEpisodeGenerationCommand extends Command
|
|
|
*
|
|
*
|
|
|
* @var string
|
|
* @var string
|
|
|
*/
|
|
*/
|
|
|
- protected $description = '批量生成分集定时任务,每5秒检查一次待处理任务,单次最多执行50秒';
|
|
|
|
|
-
|
|
|
|
|
- protected $deepSeekService;
|
|
|
|
|
-
|
|
|
|
|
- /**
|
|
|
|
|
- * Create a new command instance.
|
|
|
|
|
- *
|
|
|
|
|
- * @return void
|
|
|
|
|
- */
|
|
|
|
|
- public function __construct(DeepSeekService $deepSeekService)
|
|
|
|
|
- {
|
|
|
|
|
- parent::__construct();
|
|
|
|
|
- $this->deepSeekService = $deepSeekService;
|
|
|
|
|
- }
|
|
|
|
|
|
|
+ protected $description = '批量生成分集定时任务:并发调度不同anime_id的子进程处理,并发数默认5';
|
|
|
|
|
|
|
|
/**
|
|
/**
|
|
|
* Execute the console command.
|
|
* Execute the console command.
|
|
@@ -47,215 +34,82 @@ class ProcessBatchEpisodeGenerationCommand extends Command
|
|
|
*/
|
|
*/
|
|
|
public function handle()
|
|
public function handle()
|
|
|
{
|
|
{
|
|
|
- dLog('command')->info('开始执行批量生成分集任务...');
|
|
|
|
|
|
|
+ dLog('command')->info('开始执行批量生成分集调度任务...');
|
|
|
|
|
+
|
|
|
|
|
+ $concurrency = (int) env('BATCH_EPISODE_GENERATION_CONCURRENCY', 5);
|
|
|
|
|
+ if ($concurrency < 1) {
|
|
|
|
|
+ $concurrency = 1;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 找出所有有待处理任务的 anime_id
|
|
|
|
|
+ $animeIds = DB::table('mp_batch_episode_generation_details')
|
|
|
|
|
+ ->where('status', 'pending')
|
|
|
|
|
+ ->distinct()
|
|
|
|
|
+ ->pluck('anime_id')
|
|
|
|
|
+ ->map(function ($animeId) {
|
|
|
|
|
+ return (int) $animeId;
|
|
|
|
|
+ })
|
|
|
|
|
+ ->sort()
|
|
|
|
|
+ ->values()
|
|
|
|
|
+ ->toArray();
|
|
|
|
|
+
|
|
|
|
|
+ if (empty($animeIds)) {
|
|
|
|
|
+ dLog('command')->info('没有待处理的任务,本次调度结束');
|
|
|
|
|
+ return 0;
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
- // 每5秒执行一次任务,直到50秒后停止(参考 CheckVideoGenerationTasksCommand 的循环逻辑)
|
|
|
|
|
- $time_start = time();
|
|
|
|
|
- while (true) {
|
|
|
|
|
- $time_diff = time() - $time_start;
|
|
|
|
|
- sleep(5);
|
|
|
|
|
- if ($time_diff > 50) {
|
|
|
|
|
|
|
+ $spawned = 0;
|
|
|
|
|
+ foreach ($animeIds as $animeId) {
|
|
|
|
|
+ if ($spawned >= $concurrency) {
|
|
|
|
|
+ dLog('command')->info("已达到最大并发数 {$concurrency},停止调度");
|
|
|
break;
|
|
break;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- // 处理一个待生成的剧集;返回 1 表示处理失败,结束本次执行
|
|
|
|
|
- $result = $this->processPendingEpisodeTask();
|
|
|
|
|
- if ($result === -1) break;
|
|
|
|
|
|
|
+ $lockName = ProcessBatchAnimeEpisodeGenerationCommand::LOCK_PREFIX . $animeId;
|
|
|
|
|
+
|
|
|
|
|
+ // 仅当该 anime 的处理锁空闲时才拉起子进程;
|
|
|
|
|
+ // 真正的加锁由子进程完成,避免锁挂在调度进程的连接上
|
|
|
|
|
+ $isFree = DB::selectOne('SELECT IS_FREE_LOCK(?) AS is_free', [$lockName]);
|
|
|
|
|
+ if (!$isFree || (int) $isFree->is_free !== 1) {
|
|
|
|
|
+ dLog('command')->info("anime_id: {$animeId} 正在被其他进程处理,跳过");
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ try {
|
|
|
|
|
+ $this->spawnWorker($animeId);
|
|
|
|
|
+ $spawned++;
|
|
|
|
|
+ dLog('command')->info("已拉起子进程处理 anime_id: {$animeId}");
|
|
|
|
|
+ } catch (\Throwable $e) {
|
|
|
|
|
+ dLog('command')->error("拉起子进程失败 anime_id: {$animeId}: " . $e->getMessage());
|
|
|
|
|
+ logDB('command', 'error', '批量分集调度:拉起子进程失败', [
|
|
|
|
|
+ 'anime_id' => $animeId,
|
|
|
|
|
+ 'error' => $e->getMessage(),
|
|
|
|
|
+ 'trace' => $e->getTraceAsString()
|
|
|
|
|
+ ]);
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ dLog('command')->info("批量生成分集调度完成,本次拉起子进程数: {$spawned}");
|
|
|
return 0;
|
|
return 0;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
/**
|
|
/**
|
|
|
- * 处理一个待生成的剧集
|
|
|
|
|
|
|
+ * 拉起一个后台子进程,独占处理指定 anime_id 的待生成剧集
|
|
|
*
|
|
*
|
|
|
- * @return int 0-成功或无任务可处理(可继续下一轮),1-处理失败
|
|
|
|
|
|
|
+ * @param int $animeId
|
|
|
|
|
+ * @return void
|
|
|
*/
|
|
*/
|
|
|
- private function processPendingEpisodeTask()
|
|
|
|
|
|
|
+ private function spawnWorker($animeId)
|
|
|
{
|
|
{
|
|
|
- $anime_tasks = DB::table('mp_batch_episode_generation_details')->where('status', 'pending')->pluck('anime_id')->toArray();
|
|
|
|
|
- if (!$anime_tasks) return -1;
|
|
|
|
|
- foreach ($anime_tasks as $anime_id) {
|
|
|
|
|
- dLog('command')->info("~~~~~~开始执行($anime_id)任务~~~~~~");
|
|
|
|
|
- try {
|
|
|
|
|
- // 查找状态为 pending 的任务,按 anime_id 和 episode_number 排序
|
|
|
|
|
- $task = DB::table('mp_batch_episode_generation_details')
|
|
|
|
|
- ->where('status', 'pending')
|
|
|
|
|
- ->orderBy('anime_id')
|
|
|
|
|
- ->orderBy('episode_number')
|
|
|
|
|
- ->first();
|
|
|
|
|
-
|
|
|
|
|
- if (!$task) {
|
|
|
|
|
- dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 没有待处理的任务');
|
|
|
|
|
- continue;
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 找到待处理任务 - Anime ID: ' . $task->anime_id . ', Episode: ' . $task->episode_number);
|
|
|
|
|
-
|
|
|
|
|
- // 检查前一集是否已完成(如果不是第1集)
|
|
|
|
|
- if ($task->episode_number > 1) {
|
|
|
|
|
- $prevEpisodeNumber = $task->episode_number - 1;
|
|
|
|
|
-
|
|
|
|
|
- // 检查前一集是否在批量任务中
|
|
|
|
|
- $prevTaskInBatch = DB::table('mp_batch_episode_generation_details')
|
|
|
|
|
- ->where('anime_id', $task->anime_id)
|
|
|
|
|
- ->where('episode_number', $prevEpisodeNumber)
|
|
|
|
|
- ->first();
|
|
|
|
|
-
|
|
|
|
|
- if ($prevTaskInBatch && $prevTaskInBatch->status !== 'completed') {
|
|
|
|
|
- dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 前一集(第' . $prevEpisodeNumber . '集)尚未完成,跳过当前任务 (anime_id: ' . $task->anime_id . ')');
|
|
|
|
|
- continue;
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // 检查前一集是否已在系统中生成(查询 mp_anime_episodes)
|
|
|
|
|
- $prevEpisodeExists = DB::table('mp_anime_episodes')
|
|
|
|
|
- ->where('anime_id', $task->anime_id)
|
|
|
|
|
- ->where('episode_number', $prevEpisodeNumber)
|
|
|
|
|
- ->where('is_default', 1)
|
|
|
|
|
- ->exists();
|
|
|
|
|
-
|
|
|
|
|
- if (!$prevEpisodeExists) {
|
|
|
|
|
- dLog('command')->error('[' . date('Y-m-d H:i:s') . '] 前一集(第' . $prevEpisodeNumber . '集)不存在,无法继续生成 (anime_id: ' . $task->anime_id . ')');
|
|
|
|
|
-
|
|
|
|
|
- // 标记为失败
|
|
|
|
|
- DB::table('mp_batch_episode_generation_details')
|
|
|
|
|
- ->where('id', $task->id)
|
|
|
|
|
- ->update([
|
|
|
|
|
- 'status' => 'failed',
|
|
|
|
|
- 'error_message' => "前一集(第{$prevEpisodeNumber}集)不存在,无法继续生成",
|
|
|
|
|
- 'updated_at' => now()
|
|
|
|
|
- ]);
|
|
|
|
|
-
|
|
|
|
|
- continue;
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // 检查当前集是否已经存在(避免重复生成)
|
|
|
|
|
- $currentEpisodeExists = DB::table('mp_anime_episodes')
|
|
|
|
|
- ->where('anime_id', $task->anime_id)
|
|
|
|
|
- ->where('episode_number', $task->episode_number)
|
|
|
|
|
- ->where('is_default', 1)
|
|
|
|
|
- ->exists();
|
|
|
|
|
-
|
|
|
|
|
- if ($currentEpisodeExists) {
|
|
|
|
|
- dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 第' . $task->episode_number . '集已存在,跳过生成 (anime_id: ' . $task->anime_id . ')');
|
|
|
|
|
-
|
|
|
|
|
- // 标记为已完成
|
|
|
|
|
- DB::table('mp_batch_episode_generation_details')
|
|
|
|
|
- ->where('id', $task->id)
|
|
|
|
|
- ->update([
|
|
|
|
|
- 'status' => 'completed',
|
|
|
|
|
- 'completed_at' => now(),
|
|
|
|
|
- // 'updated_at' => now()
|
|
|
|
|
- ]);
|
|
|
|
|
-
|
|
|
|
|
- continue;
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // 标记为处理中
|
|
|
|
|
- DB::table('mp_batch_episode_generation_details')
|
|
|
|
|
- ->where('id', $task->id)
|
|
|
|
|
- ->update([
|
|
|
|
|
- 'status' => 'processing',
|
|
|
|
|
- 'updated_at' => now()
|
|
|
|
|
- ]);
|
|
|
|
|
-
|
|
|
|
|
- // 设置用户上下文(绑定到容器)
|
|
|
|
|
- app()->instance('siteData', [
|
|
|
|
|
- 'uid' => $task->uid,
|
|
|
|
|
- 'cpid' => $task->cpid
|
|
|
|
|
- ]);
|
|
|
|
|
-
|
|
|
|
|
- // 解析请求数据
|
|
|
|
|
- $requestData = json_decode($task->request_data, true);
|
|
|
|
|
- $requestData['episode_number'] = $task->episode_number;
|
|
|
|
|
-
|
|
|
|
|
- // 使用"继续策划下一集"逻辑
|
|
|
|
|
- $requestData['prompt'] = '继续策划下一集';
|
|
|
|
|
-
|
|
|
|
|
- dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 开始生成第' . $task->episode_number . '集... (anime_id: ' . $task->anime_id . ')');
|
|
|
|
|
-
|
|
|
|
|
- // 调用非流式生成方法,并记录实际执行耗时
|
|
|
|
|
- $chatStartTime = microtime(true);
|
|
|
|
|
- try {
|
|
|
|
|
- $result = $this->deepSeekService->chatForAceNonStream($requestData);
|
|
|
|
|
- } finally {
|
|
|
|
|
- $chatDuration = round(microtime(true) - $chatStartTime, 2);
|
|
|
|
|
- dLog('command')->info('[' . date('Y-m-d H:i:s') . '] chatForAceNonStream 执行完成,耗时: ' . $chatDuration . ' 秒 (anime_id: ' . $task->anime_id . ', 第' . $task->episode_number . '集)');
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // 检查是否有错误
|
|
|
|
|
- if (isset($result['error']) && $result['error']) {
|
|
|
|
|
- dLog('command')->error('[' . date('Y-m-d H:i:s') . '] 生成失败: ' . $result['error'] . ' (anime_id: ' . $task->anime_id . ', 第' . $task->episode_number . '集)');
|
|
|
|
|
-
|
|
|
|
|
- // 更新重试次数
|
|
|
|
|
- $retryCount = $task->retry_count + 1;
|
|
|
|
|
-
|
|
|
|
|
- // 标记为失败
|
|
|
|
|
- DB::table('mp_batch_episode_generation_details')
|
|
|
|
|
- ->where('id', $task->id)
|
|
|
|
|
- ->update([
|
|
|
|
|
- 'status' => 'pending',
|
|
|
|
|
- 'error_message' => $result['error'],
|
|
|
|
|
- 'retry_count' => $retryCount,
|
|
|
|
|
- 'updated_at' => now()
|
|
|
|
|
- ]);
|
|
|
|
|
-
|
|
|
|
|
- // 记录错误日志
|
|
|
|
|
- logDB('batch_episode_generation', 'error', "anime_id: {$task->anime_id} 第{$task->episode_number}集生成失败", [
|
|
|
|
|
- 'anime_id' => $task->anime_id,
|
|
|
|
|
- 'episode_number' => $task->episode_number,
|
|
|
|
|
- 'error' => $result['error']
|
|
|
|
|
- ]);
|
|
|
|
|
-
|
|
|
|
|
- continue;
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // 标记为完成
|
|
|
|
|
- DB::table('mp_batch_episode_generation_details')
|
|
|
|
|
- ->where('id', $task->id)
|
|
|
|
|
- ->update([
|
|
|
|
|
- 'status' => 'completed',
|
|
|
|
|
- 'result_data' => json_encode($result, JSON_UNESCAPED_UNICODE),
|
|
|
|
|
- 'completed_at' => now(),
|
|
|
|
|
- // 'updated_at' => now()
|
|
|
|
|
- ]);
|
|
|
|
|
-
|
|
|
|
|
- dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 第' . $task->episode_number . '集生成成功 (anime_id: ' . $task->anime_id . ')');
|
|
|
|
|
-
|
|
|
|
|
- continue;
|
|
|
|
|
-
|
|
|
|
|
- } catch (\Exception $e) {
|
|
|
|
|
- $animeLogInfo = isset($task) ? ' (anime_id: ' . $task->anime_id . ', 第' . $task->episode_number . '集)' : '';
|
|
|
|
|
- dLog('command')->error('[' . date('Y-m-d H:i:s') . '] 生成失败: ' . $e->getMessage() . $animeLogInfo);
|
|
|
|
|
-
|
|
|
|
|
- if (isset($task)) {
|
|
|
|
|
- // 更新重试次数
|
|
|
|
|
- $retryCount = $task->retry_count + 1;
|
|
|
|
|
-
|
|
|
|
|
- // 标记为失败
|
|
|
|
|
- DB::table('mp_batch_episode_generation_details')
|
|
|
|
|
- ->where('id', $task->id)
|
|
|
|
|
- ->update([
|
|
|
|
|
- 'status' => 'pending',
|
|
|
|
|
- 'error_message' => $e->getMessage(),
|
|
|
|
|
- 'retry_count' => $retryCount,
|
|
|
|
|
- 'updated_at' => now()
|
|
|
|
|
- ]);
|
|
|
|
|
-
|
|
|
|
|
- // 记录错误日志
|
|
|
|
|
- logDB('batch_episode_generation', 'error', "anime_id: {$task->anime_id} 第{$task->episode_number}集生成失败", [
|
|
|
|
|
- 'anime_id' => $task->anime_id,
|
|
|
|
|
- 'episode_number' => $task->episode_number,
|
|
|
|
|
- 'error' => $e->getMessage(),
|
|
|
|
|
- 'trace' => $e->getTraceAsString()
|
|
|
|
|
- ]);
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- continue;
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- return 0;
|
|
|
|
|
|
|
+ $process = new Process([
|
|
|
|
|
+ PHP_BINARY,
|
|
|
|
|
+ base_path('artisan'),
|
|
|
|
|
+ 'batch:generate-episodes:anime',
|
|
|
|
|
+ (string) $animeId,
|
|
|
|
|
+ ]);
|
|
|
|
|
+ $process->setTimeout(null);
|
|
|
|
|
+ $process->setIdleTimeout(null);
|
|
|
|
|
+ $process->disableOutput();
|
|
|
|
|
+ $process->start();
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|