Files
tougao/application/common/mq/ReferenceCheckArticleWorker.php

385 lines
13 KiB
PHP
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
<?php
namespace app\common\mq;
use think\Db;
use app\common\DbReconnectHelper;
use app\common\ReferenceRelevanceCheckService;
/**
* RabbitMQ 消费(队列 reference_check / ref_check.article
* 全局文章串行,文章内 reference_no 升序链式逐条「主题相关性」校对。
* 支持断点续跑:已完成条跳过,从卡死的 pending 条继续。
*/
class ReferenceCheckArticleWorker
{
const BATCH_WAITING = 0;
const BATCH_RUNNING = 1;
const BATCH_DONE = 2;
const BATCH_PARTIAL_FAILED = 3;
/** 批次/心跳超时:超过该秒数无 updated_at 更新视为僵尸,可抢占续跑 */
const BATCH_STALE_SECONDS = 1200;
/** @var ReferenceRelevanceCheckService */
private $svc;
public function __construct()
{
$this->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;
}
// 先释放其它文章上的僵尸 RUNNING避免全局串行永久堵死
$this->recoverStaleForeignBatches($batchId);
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;
}
$claim = $this->claimOrResumeBatch($batchId);
if ($claim === 'skip') {
return;
}
$resume = ($claim === 'resume');
$owned = true;
$finished = false;
try {
// 续跑时强制把卡死行收回 pending已完成行不动
$this->svc->recoverQueueRowsForArticle($pArticleId, $resume);
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
. ($resume ? ' resume=1' : '')
);
while (true) {
$row = $this->fetchNextPendingRow($pArticleId);
if (empty($row)) {
break;
}
$checkId = $this->svc->resolveCheckRowId($row);
if ($checkId <= 0) {
continue;
}
$this->processOneRow($checkId, $row, $trigger === 'recheck_pending_only');
// 每条结束后刷新批次心跳,长文不会被误判为僵尸
$this->touchBatch($batchId);
}
$stats = $this->summarizeArticleCheckStats($pArticleId);
$this->finalizeBatch($batchId, $stats['done'], $stats['failed'], $stats['total']);
$finished = true;
$this->svc->log(
'ReferenceCheckArticleWorker done p_article_id=' . $pArticleId
. ' batch_id=' . $batchId
. ' done=' . $stats['done']
. ' failed=' . $stats['failed']
);
$this->publishNextWaitingBatch();
} catch (\Throwable $e) {
// 异常不 finalize保持 RUNNING靠心跳超时后由后续消息断点续跑
if ($owned) {
$this->touchBatch($batchId);
$this->svc->log(
'ReferenceCheckArticleWorker abort batch_id=' . $batchId
. ' p_article_id=' . $pArticleId
. ' err=' . $e->getMessage()
);
}
throw $e;
} finally {
if ($owned && !$finished) {
// 消息进 DLQ 后主队列可能没人再推本批:主动再投递,便于稍后续跑
try {
(new ReferenceCheckMqPublisher())->publishArticleStart($pArticleId, $batchId, $trigger);
} catch (\Exception $pubErr) {
$this->svc->log('ReferenceCheckArticleWorker republish after abort failed: ' . $pubErr->getMessage());
}
}
}
}
/**
* @return string claim|resume|skip
*/
private function claimOrResumeBatch($batchId)
{
$batchId = intval($batchId);
$now = date('Y-m-d H:i:s');
$claimed = Db::name('article_reference_relevance_check_batch')
->where('id', $batchId)
->where('batch_status', self::BATCH_WAITING)
->update([
'batch_status' => self::BATCH_RUNNING,
'updated_at' => $now,
]);
if (intval($claimed) > 0) {
return 'claim';
}
$batch = $this->getBatch($batchId);
if (empty($batch)) {
return 'skip';
}
$status = intval($batch['batch_status']);
if ($status === self::BATCH_DONE || $status === self::BATCH_PARTIAL_FAILED) {
return 'skip';
}
if ($status === self::BATCH_RUNNING) {
if ($this->isBatchStale($batch)) {
// 僵尸 RUNNING抢占续跑不重置已完成明细
Db::name('article_reference_relevance_check_batch')
->where('id', $batchId)
->where('batch_status', self::BATCH_RUNNING)
->update(['updated_at' => $now]);
$this->svc->log('ReferenceCheckArticleWorker reclaim stale batch_id=' . $batchId);
return 'resume';
}
// 仍有心跳,说明别的消费者在跑本批
return 'skip';
}
return 'skip';
}
private function isBatchStale(array $batch)
{
$updatedAt = isset($batch['updated_at']) ? strtotime((string)$batch['updated_at']) : 0;
if ($updatedAt <= 0) {
return true;
}
return (time() - $updatedAt) >= self::BATCH_STALE_SECONDS;
}
private function touchBatch($batchId)
{
Db::name('article_reference_relevance_check_batch')
->where('id', intval($batchId))
->where('batch_status', self::BATCH_RUNNING)
->update(['updated_at' => date('Y-m-d H:i:s')]);
}
/**
* 其它文章僵尸 RUNNING → 改回 WAITING 并重新入队,从卡死条续跑
*/
private function recoverStaleForeignBatches($exceptBatchId)
{
$exceptBatchId = intval($exceptBatchId);
$staleBefore = date('Y-m-d H:i:s', time() - self::BATCH_STALE_SECONDS);
$rows = Db::name('article_reference_relevance_check_batch')
->where('batch_status', self::BATCH_RUNNING)
->where('id', '<>', $exceptBatchId)
->whereRaw('(updated_at IS NULL OR updated_at < ?)', [$staleBefore])
->order('id asc')
->limit(20)
->select();
if (empty($rows)) {
return;
}
$publisher = new ReferenceCheckMqPublisher();
foreach ($rows as $row) {
$bid = intval($row['id']);
$pid = intval($row['p_article_id']);
$affected = Db::name('article_reference_relevance_check_batch')
->where('id', $bid)
->where('batch_status', self::BATCH_RUNNING)
->whereRaw('(updated_at IS NULL OR updated_at < ?)', [$staleBefore])
->update([
'batch_status' => self::BATCH_WAITING,
'updated_at' => date('Y-m-d H:i:s'),
]);
if (intval($affected) <= 0) {
continue;
}
// 强制收回卡死行;已完成条保持不动
$this->svc->recoverQueueRowsForArticle($pid, true);
$this->svc->log('ReferenceCheckArticleWorker recover stale foreign batch_id=' . $bid . ' p_article_id=' . $pid);
try {
$publisher->publishArticleStart(
$pid,
$bid,
isset($row['trigger']) ? $row['trigger'] : 'enqueue'
);
} catch (\Exception $e) {
$this->svc->log('ReferenceCheckArticleWorker recover publish failed batch_id=' . $bid . ' err=' . $e->getMessage());
}
}
}
private function summarizeArticleCheckStats($pArticleId)
{
$rows = Db::name('article_reference_relevance_check_result')
->where('p_article_id', intval($pArticleId))
->field('status')
->select();
$done = 0;
$failed = 0;
$pending = 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++;
} elseif ($st === ReferenceRelevanceCheckService::RECORD_PENDING) {
$pending++;
}
}
return [
'done' => $done,
'failed' => $failed,
'pending' => $pending,
'total' => $done + $failed + $pending,
];
}
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 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,
'updated_at' => date('Y-m-d H:i:s'),
]);
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, $total = 0)
{
$batch = $this->getBatch($batchId);
if (empty($batch)) {
return;
}
$done = intval($done);
$failed = intval($failed);
$total = intval($total);
if ($total <= 0) {
$total = intval($batch['total_count']);
}
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());
}
}
}