svc = new ReferenceRelevanceCheckService(); } public function handleMessage(array $payload) { DbReconnectHelper::ensure(); $pArticleId = intval(isset($payload['p_article_id']) ? $payload['p_article_id'] : 0); $batchId = intval(isset($payload['batch_id']) ? $payload['batch_id'] : 0); $trigger = isset($payload['trigger']) ? (string)$payload['trigger'] : 'enqueue'; if ($pArticleId <= 0 || $batchId <= 0) { $this->svc->log('ReferenceCheckArticleWorker invalid payload'); return; } if (!$this->canStartArticleWork($batchId)) { $this->svc->log('ReferenceCheckArticleWorker defer batch_id=' . $batchId . ' other article running'); (new ReferenceCheckMqPublisher())->publishArticleStart( $pArticleId, $batchId, isset($payload['trigger']) ? $payload['trigger'] : 'enqueue' ); sleep(3); return; } if (!$this->claimBatch($batchId)) { $batch = $this->getBatch($batchId); // 已被其他消费者领取或已结束,当前消息直接跳过,避免同批次并发重复执行 if (empty($batch) || intval($batch['batch_status']) === self::BATCH_RUNNING || intval($batch['batch_status']) === self::BATCH_DONE || intval($batch['batch_status']) === self::BATCH_PARTIAL_FAILED) { return; } } $this->svc->recoverQueueRowsForArticle($pArticleId); if ($trigger !== 'recheck_pending_only' && ReferenceRelevanceCheckService::PREPARE_LITERATURE_BEFORE_CHECK) { $this->svc->prepareLiteratureContentByArticle($pArticleId); } $this->svc->log('ReferenceCheckArticleWorker start p_article_id=' . $pArticleId . ' batch_id=' . $batchId); // 快照本批待处理 id:联合引用组长一次会整组落库,循环计数会小于 total_count,收尾按快照回填 $trackedIds = $this->listPendingCheckIds($pArticleId); $done = 0; $failed = 0; while (true) { $row = $this->fetchNextPendingRow($pArticleId); if (empty($row)) { break; } $checkId = $this->svc->resolveCheckRowId($row); if ($checkId <= 0) { continue; } $result = $this->processOneRow($checkId, $row, $trigger === 'recheck_pending_only'); if ($result === 'ok') { $done++; } elseif ($result === 'failed') { $failed++; } } if (!empty($trackedIds)) { $stats = $this->summarizeTrackedCheckIds($trackedIds); $done = intval($stats['done']); $failed = intval($stats['failed']); } $this->finalizeBatch($batchId, $done, $failed); $this->svc->log('ReferenceCheckArticleWorker done p_article_id=' . $pArticleId . ' batch_id=' . $batchId . ' done=' . $done . ' failed=' . $failed); $this->publishNextWaitingBatch(); } private function listPendingCheckIds($pArticleId) { $rows = Db::name('article_reference_relevance_check_result') ->where('p_article_id', intval($pArticleId)) ->where('queue_status', ReferenceRelevanceCheckService::QUEUE_PENDING) ->where('status', ReferenceRelevanceCheckService::RECORD_PENDING) ->field('id') ->select(); $ids = []; foreach ($rows as $row) { $id = intval(isset($row['id']) ? $row['id'] : 0); if ($id > 0) { $ids[] = $id; } } return $ids; } private function summarizeTrackedCheckIds(array $checkIds) { $checkIds = array_values(array_filter(array_map('intval', $checkIds))); if (empty($checkIds)) { return ['done' => 0, 'failed' => 0]; } $rows = Db::name('article_reference_relevance_check_result') ->whereIn('id', $checkIds) ->field('id,status') ->select(); $done = 0; $failed = 0; foreach ($rows as $row) { $st = intval(isset($row['status']) ? $row['status'] : -1); if ($st === ReferenceRelevanceCheckService::RECORD_COMPLETED) { $done++; } elseif ($st === ReferenceRelevanceCheckService::RECORD_FAILED) { $failed++; } } return ['done' => $done, 'failed' => $failed]; } private function canStartArticleWork($batchId) { $running = Db::name('article_reference_relevance_check_batch') ->where('batch_status', self::BATCH_RUNNING) ->where('id', '<>', intval($batchId)) ->count(); return intval($running) === 0; } private function claimBatch($batchId) { $now = date('Y-m-d H:i:s'); $affected = Db::name('article_reference_relevance_check_batch') ->where('id', intval($batchId)) // 只允许 WAITING -> RUNNING,禁止已 RUNNING 的批次被重复 claim ->where('batch_status', self::BATCH_WAITING) ->update([ 'batch_status' => self::BATCH_RUNNING, 'updated_at' => $now, ]); return intval($affected) > 0; } private function getBatch($batchId) { return Db::name('article_reference_relevance_check_batch')->where('id', intval($batchId))->find(); } private function fetchNextPendingRow($pArticleId) { return Db::name('article_reference_relevance_check_result') ->where('p_article_id', intval($pArticleId)) ->where('queue_status', ReferenceRelevanceCheckService::QUEUE_PENDING) ->where('status', ReferenceRelevanceCheckService::RECORD_PENDING) ->order('reference_no asc,am_id asc,text_start asc,id asc') ->find(); } /** * @return string ok|failed|skip */ private function processOneRow($checkId, array $row, $skipLiteratureFetch = false) { DbReconnectHelper::ensure(); $claimed = Db::name('article_reference_relevance_check_result') ->where('id', intval($checkId)) ->where('queue_status', ReferenceRelevanceCheckService::QUEUE_PENDING) ->update(['queue_status' => ReferenceRelevanceCheckService::QUEUE_RUNNING]); if (intval($claimed) <= 0) { return 'skip'; } $retryCount = intval(isset($row['retry_count']) ? $row['retry_count'] : 0); try { $this->svc->runCheckOnce($checkId, $skipLiteratureFetch); $this->svc->markQueueRuntime($checkId, ReferenceRelevanceCheckService::QUEUE_COMPLETED, $retryCount); return 'ok'; } catch (\Exception $e) { $this->svc->log('ReferenceCheckArticleWorker check_id=' . $checkId . ' err=' . $e->getMessage()); DbReconnectHelper::ensure(); try { $fresh = Db::name('article_reference_relevance_check_result')->where('id', intval($checkId))->find(); if (!empty($fresh) && intval($fresh['status']) === ReferenceRelevanceCheckService::RECORD_FAILED) { if (intval($fresh['queue_status']) !== ReferenceRelevanceCheckService::QUEUE_FAILED) { $this->svc->markQueueRuntime($checkId, ReferenceRelevanceCheckService::QUEUE_FAILED, $retryCount); } return 'failed'; } $groupRows = !empty($fresh) ? $this->svc->findCitationGroupRowsForWorker($fresh) : []; if (!empty($groupRows)) { $this->svc->failGroupWithQueue($groupRows, $e->getMessage(), $retryCount); } else { $this->svc->updateCheckResult($checkId, [ 'status' => ReferenceRelevanceCheckService::RECORD_FAILED, 'error_msg' => $e->getMessage(), ]); $this->svc->markQueueRuntime($checkId, ReferenceRelevanceCheckService::QUEUE_FAILED, $retryCount); } } catch (\Exception $e2) { \think\Log::error('ReferenceCheckArticleWorker markFailed: ' . $e2->getMessage()); } return 'failed'; } } private function finalizeBatch($batchId, $done, $failed) { $batch = $this->getBatch($batchId); if (empty($batch)) { return; } $total = intval($batch['total_count']); $done = intval($done); $failed = intval($failed); // 快照回填后若实际终态条数多于入队 total,抬升 total 保持一致 if (($done + $failed) > $total) { $total = $done + $failed; } $status = self::BATCH_DONE; if ($failed > 0) { $status = self::BATCH_PARTIAL_FAILED; } Db::name('article_reference_relevance_check_batch')->where('id', intval($batchId))->update([ 'batch_status' => $status, 'total_count' => $total, 'done_count' => $done, 'failed_count' => $failed, 'updated_at' => date('Y-m-d H:i:s'), ]); if ($total > 0 && ($done + $failed) < $total) { $this->svc->log('ReferenceCheckArticleWorker batch_id=' . $batchId . ' incomplete total=' . $total . ' done=' . $done . ' failed=' . $failed); } } private function publishNextWaitingBatch() { $next = Db::name('article_reference_relevance_check_batch') ->where('batch_status', self::BATCH_WAITING) ->order('id asc') ->find(); if (empty($next)) { return; } try { (new ReferenceCheckMqPublisher())->publishArticleStart( intval($next['p_article_id']), intval($next['id']), isset($next['trigger']) ? $next['trigger'] : 'enqueue' ); } catch (\Exception $e) { $this->svc->log('ReferenceCheck publishNextWaitingBatch failed: ' . $e->getMessage()); \think\Log::error('ReferenceCheck publishNextWaitingBatch: ' . $e->getMessage()); } } }