fix(dispatcher): 僵尸PROCESSING任务回收 + fam-core文件日志
- db_layer 新增 reclaim_stale_processing: PROCESSING 超过 push_timeout+120s 重置 PENDING - dispatcher 轮询前先回收僵尸任务(进程重启/Edge重启导致 in-flight 请求丢失的场景) - logger 增加 fam-core/logs/fam-core.log 文件输出(daemon 模式 stdout 不可见)
This commit is contained in:
@@ -118,6 +118,34 @@ def increment_retry(task_id: int, next_retry_at: datetime):
|
||||
conn.close()
|
||||
|
||||
|
||||
def reclaim_stale_processing(timeout_seconds: int) -> List[int]:
|
||||
"""回收僵尸 PROCESSING 任务:updated_at 早于 timeout_seconds 前的任务重置为 PENDING
|
||||
|
||||
场景:Dispatcher 推送过程中进程重启/Edge 重启导致 in-flight 请求丢失,
|
||||
任务停留在 PROCESSING 无人处理。重置后由常规重试机制接管。
|
||||
返回被回收的 task_id 列表。
|
||||
"""
|
||||
conn = get_conn()
|
||||
try:
|
||||
cursor = conn.cursor()
|
||||
cursor.execute(
|
||||
"SELECT task_id FROM process_tasks "
|
||||
"WHERE status='PROCESSING' AND updated_at < NOW() - INTERVAL %s SECOND",
|
||||
(timeout_seconds,)
|
||||
)
|
||||
task_ids = [row[0] for row in cursor.fetchall()]
|
||||
if task_ids:
|
||||
placeholders = ','.join(['%s'] * len(task_ids))
|
||||
cursor.execute(
|
||||
f"UPDATE process_tasks SET status='PENDING' WHERE task_id IN ({placeholders})",
|
||||
task_ids
|
||||
)
|
||||
conn.commit()
|
||||
return task_ids
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def get_task(task_id: int) -> Optional[Dict]:
|
||||
"""获取单个任务"""
|
||||
conn = get_conn()
|
||||
|
||||
@@ -149,6 +149,16 @@ class Dispatcher:
|
||||
|
||||
def _poll_once(self):
|
||||
"""执行一次轮询"""
|
||||
# 先回收僵尸任务:PROCESSING 超过 push_timeout+缓冲 说明推送进程已丢失
|
||||
# (如 fam-core 重启、Edge 重启掐断 in-flight 连接),重置回 PENDING 走重试
|
||||
stale_timeout = self.push_timeout + 120
|
||||
try:
|
||||
stale_ids = db_layer.reclaim_stale_processing(stale_timeout)
|
||||
for tid in stale_ids:
|
||||
logger.warning(f"[task_id={tid}] PROCESSING 超时 {stale_timeout}s,回收为 PENDING 重试")
|
||||
except Exception as e:
|
||||
logger.error(f"僵尸任务回收失败: {e}", exc_info=True)
|
||||
|
||||
tasks = db_layer.get_pending_tasks(limit=10)
|
||||
for task in tasks:
|
||||
if self._should_retry(task):
|
||||
|
||||
@@ -2,23 +2,36 @@
|
||||
日志工具 - 统一格式,带 task_id 作为 trace_id
|
||||
"""
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
from datetime import datetime
|
||||
|
||||
# 日志目录:fam-core/logs/(相对 src 的上一级),失败则退化为仅 stdout
|
||||
_LOG_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), 'logs')
|
||||
|
||||
|
||||
def setup_logger(name='fam-core', level=logging.INFO):
|
||||
"""配置并返回 logger"""
|
||||
"""配置并返回 logger(stdout + 文件双写)"""
|
||||
logger = logging.getLogger(name)
|
||||
if logger.handlers:
|
||||
return logger
|
||||
logger.setLevel(level)
|
||||
handler = logging.StreamHandler(sys.stdout)
|
||||
formatter = logging.Formatter(
|
||||
'%(asctime)s [%(name)s] %(levelname)s %(message)s',
|
||||
datefmt='%Y-%m-%d %H:%M:%S'
|
||||
)
|
||||
handler = logging.StreamHandler(sys.stdout)
|
||||
handler.setFormatter(formatter)
|
||||
logger.addHandler(handler)
|
||||
# 文件输出(daemon 模式下 stdout 不可见,文件是唯一可追溯日志)
|
||||
try:
|
||||
os.makedirs(_LOG_DIR, exist_ok=True)
|
||||
file_handler = logging.FileHandler(
|
||||
os.path.join(_LOG_DIR, 'fam-core.log'), encoding='utf-8')
|
||||
file_handler.setFormatter(formatter)
|
||||
logger.addHandler(file_handler)
|
||||
except OSError:
|
||||
pass # 目录不可写时退化为仅 stdout
|
||||
return logger
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user