ProcessBatchAnimeEpisodeGenerationCommand.php 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334
  1. <?php
  2. namespace App\Console\Commands;
  3. use Illuminate\Console\Command;
  4. use Illuminate\Support\Facades\DB;
  5. use App\Services\DeepSeek\DeepSeekService;
  6. /**
  7. * 批量生成分集子进程 worker
  8. * 由 batch:generate-episodes 调度,每个进程通过 MySQL 命名锁独占一个 anime_id,
  9. * 串行生成该 anime 的 pending 剧集(按集数顺序),直到没有待处理任务为止。
  10. */
  11. class ProcessBatchAnimeEpisodeGenerationCommand extends Command
  12. {
  13. /**
  14. * anime 处理锁前缀(配合 MySQL GET_LOCK 使用,同一 anime 同一时刻只允许一个进程处理)
  15. */
  16. public const LOCK_PREFIX = 'batch_episode_generation_anime_';
  17. /**
  18. * 单集最大重试次数,超过后停止重试并标记为失败
  19. */
  20. private const MAX_RETRY_COUNT = 5;
  21. /**
  22. * The name and signature of the console command.
  23. *
  24. * @var string
  25. */
  26. protected $signature = 'batch:generate-episodes:anime {anime_id : 动漫ID}';
  27. /**
  28. * The console command description.
  29. *
  30. * @var string
  31. */
  32. protected $description = '批量生成分集子进程:独占处理指定anime_id的待生成剧集';
  33. protected $deepSeekService;
  34. /**
  35. * Create a new command instance.
  36. *
  37. * @return void
  38. */
  39. public function __construct(DeepSeekService $deepSeekService)
  40. {
  41. parent::__construct();
  42. $this->deepSeekService = $deepSeekService;
  43. }
  44. /**
  45. * Execute the console command.
  46. *
  47. * @return int
  48. */
  49. public function handle()
  50. {
  51. set_time_limit(0);
  52. $animeId = (int) $this->argument('anime_id');
  53. $lockName = self::LOCK_PREFIX . $animeId;
  54. // 非阻塞抢占该 anime 的处理锁;抢不到说明已有其他进程正在处理
  55. $acquired = DB::selectOne('SELECT GET_LOCK(?, 0) AS acquired', [$lockName]);
  56. if (!$acquired || (int) $acquired->acquired !== 1) {
  57. dLog('command')->info("anime_id: {$animeId} 正在被其他进程处理,本次直接退出");
  58. return 0;
  59. }
  60. dLog('command')->info("anime_id: {$animeId} 抢占处理锁成功,开始处理...");
  61. try {
  62. $this->processAnimeEpisodes($animeId);
  63. } catch (\Throwable $e) {
  64. dLog('command')->error("anime_id: {$animeId} 处理异常: " . $e->getMessage());
  65. logDB('batch_episode_generation', 'error', "anime_id: {$animeId} 批量生成处理异常", [
  66. 'anime_id' => $animeId,
  67. 'error' => $e->getMessage(),
  68. 'trace' => $e->getTraceAsString()
  69. ]);
  70. } finally {
  71. // 显式释放处理锁;即使进程异常退出,连接断开后 MySQL 也会自动释放
  72. DB::statement('SELECT RELEASE_LOCK(?)', [$lockName]);
  73. dLog('command')->info("anime_id: {$animeId} 处理结束,已释放处理锁");
  74. }
  75. return 0;
  76. }
  77. /**
  78. * 串行处理指定 anime 的所有待生成剧集
  79. *
  80. * @param int $animeId
  81. * @return void
  82. */
  83. private function processAnimeEpisodes($animeId)
  84. {
  85. while (true) {
  86. // 取该 anime 集数最小的 pending 任务
  87. $task = DB::table('mp_batch_episode_generation_details')
  88. ->where('anime_id', $animeId)
  89. ->where('status', 'pending')
  90. ->orderBy('episode_number')
  91. ->first();
  92. if (!$task) {
  93. dLog('command')->info('[' . date('Y-m-d H:i:s') . "] anime_id: {$animeId} 没有待处理的任务");
  94. return;
  95. }
  96. try {
  97. // 返回 false 表示该 anime 暂时无法继续(如前一集未完成),退出等待下一轮调度
  98. $keepGoing = $this->processEpisode($task);
  99. if (!$keepGoing) {
  100. dLog('command')->info("anime_id: {$animeId} 暂时无法继续生成,退出本次处理");
  101. return;
  102. }
  103. } catch (\Exception $e) {
  104. $animeLogInfo = isset($task) ? ' (anime_id: ' . $task->anime_id . ', 第' . $task->episode_number . '集)' : '';
  105. dLog('command')->error('[' . date('Y-m-d H:i:s') . '] 生成失败: ' . $e->getMessage() . $animeLogInfo);
  106. if (isset($task)) {
  107. // 更新重试次数
  108. $retryCount = $task->retry_count + 1;
  109. // 未超过最大重试次数则标记为待处理继续重试,超过后标记为失败并通知
  110. $this->markTaskRetryOrFail($task, $retryCount, $e->getMessage());
  111. // 记录错误日志
  112. logDB('batch_episode_generation', 'error', "anime_id: {$task->anime_id} 第{$task->episode_number}集生成失败", [
  113. 'anime_id' => $task->anime_id,
  114. 'episode_number' => $task->episode_number,
  115. 'error' => $e->getMessage(),
  116. 'trace' => $e->getTraceAsString()
  117. ]);
  118. }
  119. }
  120. // 与原有逻辑保持一致:每5秒处理一轮
  121. sleep(5);
  122. }
  123. }
  124. /**
  125. * 处理一个待生成的剧集(沿用原有校验与生成逻辑)
  126. *
  127. * @param object $task
  128. * @return bool true-可继续处理下一个任务,false-该anime暂时无法继续
  129. */
  130. private function processEpisode($task)
  131. {
  132. dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 找到待处理任务 - Anime ID: ' . $task->anime_id . ', Episode: ' . $task->episode_number);
  133. // 检查前一集是否已完成(如果不是第1集)
  134. if ($task->episode_number > 1) {
  135. $prevEpisodeNumber = $task->episode_number - 1;
  136. // 优先判断前一集是否已在系统中生成(查询 mp_anime_episodes)。
  137. // 用户手动创建的前一集也视为已存在,直接跳过批量任务状态判断,
  138. // 避免批量任务前一集失败后卡住整个批量任务。
  139. $prevEpisodeExists = DB::table('mp_anime_episodes')
  140. ->where('anime_id', $task->anime_id)
  141. ->where('episode_number', $prevEpisodeNumber)
  142. ->where('is_default', 1)
  143. ->exists();
  144. if (!$prevEpisodeExists) {
  145. // 系统中不存在前一集时,才进行批量任务状态判断
  146. $prevTaskInBatch = DB::table('mp_batch_episode_generation_details')
  147. ->where('anime_id', $task->anime_id)
  148. ->where('episode_number', $prevEpisodeNumber)
  149. ->first();
  150. if ($prevTaskInBatch && $prevTaskInBatch->status !== 'completed') {
  151. dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 前一集(第' . $prevEpisodeNumber . '集)尚未完成,跳过当前任务 (anime_id: ' . $task->anime_id . ')');
  152. return false;
  153. }
  154. dLog('command')->error('[' . date('Y-m-d H:i:s') . '] 前一集(第' . $prevEpisodeNumber . '集)不存在,无法继续生成 (anime_id: ' . $task->anime_id . ')');
  155. // 标记为失败
  156. DB::table('mp_batch_episode_generation_details')
  157. ->where('id', $task->id)
  158. ->update([
  159. 'status' => 'failed',
  160. 'error_message' => "前一集(第{$prevEpisodeNumber}集)不存在,无法继续生成",
  161. 'updated_at' => now()
  162. ]);
  163. return false;
  164. }
  165. }
  166. // 检查当前集是否已经存在(避免重复生成)
  167. $currentEpisodeExists = DB::table('mp_anime_episodes')
  168. ->where('anime_id', $task->anime_id)
  169. ->where('episode_number', $task->episode_number)
  170. ->where('is_default', 1)
  171. ->exists();
  172. if ($currentEpisodeExists) {
  173. dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 第' . $task->episode_number . '集已存在,跳过生成 (anime_id: ' . $task->anime_id . ')');
  174. // 标记为已完成
  175. DB::table('mp_batch_episode_generation_details')
  176. ->where('id', $task->id)
  177. ->update([
  178. 'status' => 'completed',
  179. 'completed_at' => now(),
  180. // 'updated_at' => now()
  181. ]);
  182. return true;
  183. }
  184. // 标记为处理中
  185. DB::table('mp_batch_episode_generation_details')
  186. ->where('id', $task->id)
  187. ->update([
  188. 'status' => 'processing',
  189. 'updated_at' => now()
  190. ]);
  191. // 设置用户上下文(绑定到容器)
  192. app()->instance('siteData', [
  193. 'uid' => $task->uid,
  194. 'cpid' => $task->cpid
  195. ]);
  196. // 解析请求数据
  197. $requestData = json_decode($task->request_data, true);
  198. $requestData['episode_number'] = $task->episode_number;
  199. // 使用"继续策划下一集"逻辑
  200. $requestData['prompt'] = '继续策划下一集';
  201. dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 开始生成第' . $task->episode_number . '集... (anime_id: ' . $task->anime_id . ')');
  202. // 调用非流式生成方法,并记录实际执行耗时
  203. $chatStartTime = microtime(true);
  204. try {
  205. $result = $this->deepSeekService->chatForAceNonStream($requestData);
  206. } finally {
  207. $chatDuration = round(microtime(true) - $chatStartTime, 2);
  208. dLog('command')->info('[' . date('Y-m-d H:i:s') . '] chatForAceNonStream 执行完成,耗时: ' . $chatDuration . ' 秒 (anime_id: ' . $task->anime_id . ', 第' . $task->episode_number . '集)');
  209. }
  210. // 检查是否有错误
  211. if (isset($result['error']) && $result['error']) {
  212. dLog('command')->error('[' . date('Y-m-d H:i:s') . '] 生成失败: ' . $result['error'] . ' (anime_id: ' . $task->anime_id . ', 第' . $task->episode_number . '集)');
  213. // 更新重试次数
  214. $retryCount = $task->retry_count + 1;
  215. // 未超过最大重试次数则标记为待处理继续重试,超过后标记为失败并通知
  216. $this->markTaskRetryOrFail($task, $retryCount, $result['error']);
  217. // 记录错误日志
  218. logDB('batch_episode_generation', 'error', "anime_id: {$task->anime_id} 第{$task->episode_number}集生成失败", [
  219. 'anime_id' => $task->anime_id,
  220. 'episode_number' => $task->episode_number,
  221. 'error' => $result['error']
  222. ]);
  223. return true;
  224. }
  225. // 标记为完成
  226. DB::table('mp_batch_episode_generation_details')
  227. ->where('id', $task->id)
  228. ->update([
  229. 'status' => 'completed',
  230. 'result_data' => json_encode($result, JSON_UNESCAPED_UNICODE),
  231. 'completed_at' => now(),
  232. // 'updated_at' => now()
  233. ]);
  234. dLog('command')->info('[' . date('Y-m-d H:i:s') . '] 第' . $task->episode_number . '集生成成功 (anime_id: ' . $task->anime_id . ')');
  235. return true;
  236. }
  237. /**
  238. * 任务失败处理:未超过最大重试次数则标记为待处理继续重试,
  239. * 超过后直接标记为失败并通过 sendNotice 发送报错通知。
  240. *
  241. * @param object $task
  242. * @param int $retryCount
  243. * @param string $error
  244. * @return void
  245. */
  246. private function markTaskRetryOrFail($task, $retryCount, $error)
  247. {
  248. if ($retryCount > self::MAX_RETRY_COUNT) {
  249. // 超过最大重试次数,标记为失败
  250. DB::table('mp_batch_episode_generation_details')
  251. ->where('id', $task->id)
  252. ->update([
  253. 'status' => 'failed',
  254. 'error_message' => '重试超过' . self::MAX_RETRY_COUNT . '次,已停止重试: ' . $error,
  255. 'retry_count' => $retryCount,
  256. 'updated_at' => now()
  257. ]);
  258. dLog('command')->error('[' . date('Y-m-d H:i:s') . '] 第' . $task->episode_number . '集重试超过' . self::MAX_RETRY_COUNT . '次,已标记为失败 (anime_id: ' . $task->anime_id . ')');
  259. // 发送报错通知(通知失败不影响任务状态)
  260. try {
  261. $notice = "批量分集生成失败(重试超过" . self::MAX_RETRY_COUNT . "次,已停止重试)\n"
  262. . "anime_id: {$task->anime_id}\n"
  263. . "集数: 第{$task->episode_number}集\n"
  264. . "错误: {$error}\n"
  265. . "时间: " . date('Y-m-d H:i:s');
  266. sendNotice($notice);
  267. } catch (\Throwable $e) {
  268. dLog('command')->error('[' . date('Y-m-d H:i:s') . '] 发送失败通知异常 (anime_id: ' . $task->anime_id . ', 第' . $task->episode_number . '集): ' . $e->getMessage());
  269. }
  270. return;
  271. }
  272. // 未超过最大重试次数,标记为待处理,下一轮重试
  273. DB::table('mp_batch_episode_generation_details')
  274. ->where('id', $task->id)
  275. ->update([
  276. 'status' => 'pending',
  277. 'error_message' => $error,
  278. 'retry_count' => $retryCount,
  279. 'updated_at' => now()
  280. ]);
  281. }
  282. }