From 900cf57bb34cbf56d6b77f575ad118a2d2915ca7 Mon Sep 17 00:00:00 2001 From: zk Date: Thu, 9 Jul 2026 15:30:02 +0800 Subject: [PATCH] =?UTF-8?q?=E6=96=B9=E6=A1=88=E4=BF=AE=E6=94=B9=E4=B8=BA?= =?UTF-8?q?=20=E5=A4=9A=E7=BA=BF=E7=A8=8B=20=E6=BB=9A=E5=8A=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env | 10 +-- .env.prod | 10 +-- .env.test | 10 +-- app/config/settings.py | 10 +-- app/main.py | 33 +++++--- app/scheduler/tasks.py | 45 ++-------- app/services/company_clean_service.py | 78 ++++++++++------- app/services/job_clean_service.py | 116 +++++++++++++++----------- 8 files changed, 160 insertions(+), 152 deletions(-) diff --git a/.env b/.env index 82f976e..9bd33f5 100644 --- a/.env +++ b/.env @@ -24,15 +24,13 @@ ANTHROPIC_API_KEY=sk-43ccdb29caa7e9ebe0db8ac0958c63f6d3a2d62e59064d3d26d94332055 ANTHROPIC_BASE_URL=https://code.warpdevloper.cloud # 岗位清洗参数 -CLEAN_BATCH_SIZE=100 -CLEAN_CONCURRENCY=60 -CLEAN_INTERVAL_SECONDS=200 +CLEAN_WORKER_COUNT=60 +CLEAN_IDLE_SLEEP=5.0 CLEAN_TOTAL_LIMIT=0 # 公司补充参数 -COMPANY_BATCH_SIZE=20 -COMPANY_CONCURRENCY=10 -COMPANY_INTERVAL_SECONDS=300 +COMPANY_WORKER_COUNT=10 +COMPANY_IDLE_SLEEP=10.0 # 岗位下架 JOB_EXPIRE_DAYS=7 diff --git a/.env.prod b/.env.prod index e7418b8..eaeff5c 100644 --- a/.env.prod +++ b/.env.prod @@ -24,15 +24,13 @@ ANTHROPIC_API_KEY=sk-43ccdb29caa7e9ebe0db8ac0958c63f6d3a2d62e59064d3d26d94332055 ANTHROPIC_BASE_URL=https://code.warpdevloper.cloud # 岗位清洗参数 -CLEAN_BATCH_SIZE=100 -CLEAN_CONCURRENCY=60 -CLEAN_INTERVAL_SECONDS=200 +CLEAN_WORKER_COUNT=60 +CLEAN_IDLE_SLEEP=5.0 CLEAN_TOTAL_LIMIT=0 # 公司补充参数 -COMPANY_BATCH_SIZE=20 -COMPANY_CONCURRENCY=10 -COMPANY_INTERVAL_SECONDS=300 +COMPANY_WORKER_COUNT=10 +COMPANY_IDLE_SLEEP=10.0 # 岗位下架 JOB_EXPIRE_DAYS=7 diff --git a/.env.test b/.env.test index bf748c6..d698a10 100644 --- a/.env.test +++ b/.env.test @@ -20,14 +20,12 @@ VOLCENGINE_API_KEY=fd065993-bee2-4f31-8bf2-56d5d3012c02 VOLCENGINE_BASE_URL=https://ark.cn-beijing.volces.com/api/v3 # 岗位清洗参数 -CLEAN_BATCH_SIZE=100 -CLEAN_CONCURRENCY=50 -CLEAN_INTERVAL_SECONDS=180 +CLEAN_WORKER_COUNT=60 +CLEAN_IDLE_SLEEP=5.0 # 公司补充参数 -COMPANY_BATCH_SIZE=20 -COMPANY_CONCURRENCY=10 -COMPANY_INTERVAL_SECONDS=300 +COMPANY_WORKER_COUNT=10 +COMPANY_IDLE_SLEEP=10.0 # 岗位下架 JOB_EXPIRE_DAYS=7 diff --git a/app/config/settings.py b/app/config/settings.py index 719d3f5..e90e9b2 100644 --- a/app/config/settings.py +++ b/app/config/settings.py @@ -35,15 +35,13 @@ class Settings(BaseSettings): anthropic_base_url: str = "https://code.warpdevloper.cloud" # ──────────── 岗位清洗参数 ──────────── - clean_batch_size: int = 100 - clean_concurrency: int = 80 - clean_interval_seconds: int = 200 + clean_worker_count: int = 60 # 持续消费 worker 协程数量 + clean_idle_sleep: float = 5.0 # 没数据时每个 worker 休眠秒数 clean_total_limit: int = 0 # 累计清洗总数上限,达到后停止清洗任务;0 = 不限制 # ──────────── 公司补充参数 ──────────── - company_batch_size: int = 20 - company_concurrency: int = 10 - company_interval_seconds: int = 300 + company_worker_count: int = 10 # 持续消费 worker 协程数量 + company_idle_sleep: float = 10.0 # 没数据时每个 worker 休眠秒数 # ──────────── 岗位下架参数 ──────────── job_expire_days: int = 7 diff --git a/app/main.py b/app/main.py index 7653174..741e792 100644 --- a/app/main.py +++ b/app/main.py @@ -1,16 +1,16 @@ -"""项目入口:初始化数据源、加载字典、启动调度器""" +"""项目入口:初始化数据源、加载字典、启动持续消费 worker""" import asyncio import logging import signal import warnings -from datetime import datetime from app.core.logger import log # 屏蔽 asyncmy INSERT IGNORE 产生的 Duplicate entry warnings warnings.filterwarnings("ignore", message=".*Duplicate entry.*") logging.getLogger("asyncmy").setLevel(logging.ERROR) + from app.core.database import init_db, close_db from app.services.dict_cache_service import dict_cache from app.scheduler.tasks import create_scheduler @@ -27,25 +27,29 @@ async def main(): # 加载字典缓存 await dict_cache.refresh() - # 创建并启动调度器 + # 创建并启动调度器(僵尸恢复、岗位下架等辅助任务) scheduler = create_scheduler() scheduler.start() + log.info("调度器已启动,辅助定时任务已注册") - # 立即触发一次岗位清洗和公司补充 - scheduler.modify_job("job_clean", next_run_time=datetime.now()) - scheduler.modify_job("company_clean", next_run_time=datetime.now()) + # 启动持续消费任务 + from app.services.job_clean_service import run_job_clean, stop_job_clean + from app.services.company_clean_service import run_company_clean, stop_company_clean - log.info("调度器已启动,所有定时任务已注册") + job_clean_task = asyncio.create_task(run_job_clean()) + company_clean_task = asyncio.create_task(run_company_clean()) + log.info("岗位清洗 & 公司补充持续消费任务已启动") # 优雅关闭 stop_event = asyncio.Event() def _shutdown(*args): log.info("收到关闭信号,正在关闭...") + stop_job_clean() + stop_company_clean() stop_event.set() loop = asyncio.get_running_loop() - # Unix: SIGINT + SIGTERM,Windows: 仅靠 KeyboardInterrupt for sig in (signal.SIGINT, signal.SIGTERM): try: loop.add_signal_handler(sig, _shutdown) @@ -55,11 +59,14 @@ async def main(): try: await stop_event.wait() except KeyboardInterrupt: - pass - finally: - scheduler.shutdown(wait=False) - await close_db() - log.info("OfferPie Job Cleaner 已关闭") + _shutdown() + + # 等待 worker 协程退出 + await asyncio.gather(job_clean_task, company_clean_task, return_exceptions=True) + + scheduler.shutdown(wait=False) + await close_db() + log.info("OfferPie Job Cleaner 已关闭") if __name__ == "__main__": diff --git a/app/scheduler/tasks.py b/app/scheduler/tasks.py index 621da1e..05f2976 100644 --- a/app/scheduler/tasks.py +++ b/app/scheduler/tasks.py @@ -1,6 +1,11 @@ -"""定时任务注册""" +"""定时任务注册 -from datetime import datetime, timedelta +岗位清洗和公司补充改为持续消费模型(启动即运行), +其余辅助任务保留定时调度。 +""" + +import asyncio +from datetime import datetime from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.interval import IntervalTrigger @@ -11,30 +16,12 @@ from app.core.logger import log def create_scheduler() -> AsyncIOScheduler: - """创建并注册所有定时任务""" + """创建并注册辅助定时任务(岗位清洗/公司补充由 main 直接启动)""" scheduler = AsyncIOScheduler( timezone="Asia/Shanghai", job_defaults={"misfire_grace_time": 60}, ) - # 岗位清洗(每 N 秒) - scheduler.add_job( - _job_clean_task, - trigger=IntervalTrigger(seconds=settings.clean_interval_seconds), - id="job_clean", - name="岗位清洗", - max_instances=1, - ) - - # 公司补充(每 N 秒) - scheduler.add_job( - _company_clean_task, - trigger=IntervalTrigger(seconds=settings.company_interval_seconds), - id="company_clean", - name="公司补充", - max_instances=1, - ) - # 岗位僵尸恢复(每30分钟) scheduler.add_job( _job_zombie_task, @@ -65,22 +52,6 @@ def create_scheduler() -> AsyncIOScheduler: return scheduler -async def _job_clean_task(): - from app.services.job_clean_service import run_job_clean - try: - await run_job_clean() - except Exception as e: - log.error("岗位清洗任务异常: {}", e) - - -async def _company_clean_task(): - from app.services.company_clean_service import run_company_clean - try: - await run_company_clean() - except Exception as e: - log.error("公司补充任务异常: {}", e) - - async def _job_zombie_task(): from app.services.zombie_recover_service import recover_job_zombie try: diff --git a/app/services/company_clean_service.py b/app/services/company_clean_service.py index 0d2de32..768b418 100644 --- a/app/services/company_clean_service.py +++ b/app/services/company_clean_service.py @@ -1,4 +1,8 @@ -"""公司数据补充服务(协程版)""" +"""公司数据补充服务(持续消费模型) + +启动 N 个 worker 协程,每个 worker 循环:取一条 → 处理 → 取下一条。 +没数据时短暂休眠后重试。 +""" import asyncio from datetime import datetime @@ -13,47 +17,65 @@ from app.ai.prompts import COMPANY_ENRICH_SYSTEM from app.services.ai_tool import ai_chat_json from app.services.dict_cache_service import dict_cache +# 停止信号 +_stop_event = asyncio.Event() + + +def stop_company_clean(): + """外部调用,通知所有 worker 停止""" + _stop_event.set() + async def run_company_clean() -> None: - """一次批量公司补充任务""" - # 锁定一批待完善公司 + """启动 N 个 worker 协程持续消费""" + _stop_event.clear() + worker_count = settings.company_worker_count + log.info("公司补充:启动 {} 个 worker 协程", worker_count) + + workers = [asyncio.create_task(_worker(i)) for i in range(worker_count)] + await asyncio.gather(*workers) + + log.info("公司补充:所有 worker 已退出") + + +async def _worker(worker_id: int) -> None: + """单个 worker:循环取一条、处理一条""" + while not _stop_event.is_set(): + data = await _fetch_one() + if data is None: + await asyncio.sleep(settings.company_idle_sleep) + continue + + try: + await _do_clean(data) + except Exception as e: + log.error("[worker-{}] 公司补充异常, id={}, shortName={}: {}", + worker_id, data["id"], data.get("short_name"), e) + + +async def _fetch_one() -> dict | None: + """从 MySQL 锁定一条待完善公司并标记为 status=3""" async with MysqlSession() as mysql: + # MySQL 不支持 UPDATE ... RETURNING,分两步 result = await mysql.execute( text(""" SELECT * FROM bg_company WHERE status = 0 - LIMIT :limit + LIMIT 1 FOR UPDATE SKIP LOCKED """), - {"limit": settings.company_batch_size}, ) - rows = result.mappings().all() - if not rows: - return + row = result.mappings().first() + if not row: + return None - ids = [r["id"] for r in rows] - # MySQL 批量 IN 用 format 拼接(id 是 bigint,安全) - ids_str = ",".join(str(i) for i in ids) + company_id = row["id"] await mysql.execute( - text(f"UPDATE bg_company SET status = 3, update_time = NOW() WHERE id IN ({ids_str})"), + text("UPDATE bg_company SET status = 3, update_time = NOW() WHERE id = :id"), + {"id": company_id}, ) await mysql.commit() - - log.info("公司补充:锁定{}条数据", len(rows)) - - # 协程并发,信号量限流 - sem = asyncio.Semaphore(settings.company_concurrency) - tasks = [_clean_one(sem, dict(r)) for r in rows] - await asyncio.gather(*tasks, return_exceptions=True) - - -async def _clean_one(sem: asyncio.Semaphore, company: dict) -> None: - """单条公司补充""" - async with sem: - try: - await _do_clean(company) - except Exception as e: - log.error("公司补充异常, id={}, shortName={}: {}", company["id"], company.get("short_name"), e) + return dict(row) async def _do_clean(company: dict) -> None: diff --git a/app/services/job_clean_service.py b/app/services/job_clean_service.py index 888ed11..5a0e4d9 100644 --- a/app/services/job_clean_service.py +++ b/app/services/job_clean_service.py @@ -1,4 +1,8 @@ -"""岗位清洗服务(协程版)""" +"""岗位清洗服务(持续消费模型) + +启动 N 个 worker 协程,每个 worker 循环:取一条 → 处理 → 取下一条。 +没数据时短暂休眠后重试。 +""" import asyncio import json @@ -28,70 +32,83 @@ _company_lock = asyncio.Lock() # 累计已清洗总数(用于总数上限控制,进程内统计) _cleaned_total = 0 +# 停止信号 +_stop_event = asyncio.Event() + def is_clean_limit_reached() -> bool: """是否已达到累计清洗总数上限(clean_total_limit=0 表示不限制)""" return 0 < settings.clean_total_limit <= _cleaned_total +def stop_job_clean(): + """外部调用,通知所有 worker 停止""" + _stop_event.set() + + async def run_job_clean() -> None: - """一次批量清洗任务""" + """启动 N 个 worker 协程持续消费,直到收到停止信号或达到上限""" + _stop_event.clear() + worker_count = settings.clean_worker_count + log.info("岗位清洗:启动 {} 个 worker 协程", worker_count) + + workers = [asyncio.create_task(_worker(i)) for i in range(worker_count)] + await asyncio.gather(*workers) + + log.info("岗位清洗:所有 worker 已退出,累计清洗 {} 条", _cleaned_total) + + +async def _worker(worker_id: int) -> None: + """单个 worker:循环取一条、处理一条""" global _cleaned_total - # 总数上限:已达上限直接跳过 - if is_clean_limit_reached(): - return + while not _stop_event.is_set(): + # 总数上限检查 + if is_clean_limit_reached(): + return - # 1. 从 PG 锁定一批待清洗数据 + # 从 PG 取一条待处理数据 + data = await _fetch_one() + if data is None: + # 没数据,休眠后重试 + await asyncio.sleep(settings.clean_idle_sleep) + continue + + # 处理 + try: + await _do_clean(data) + _cleaned_total += 1 + except Exception as e: + log.error("[worker-{}] 岗位清洗异常, id={}: {}", worker_id, data["id"], e) + + # 达到上限后通知所有 worker 停止 + if is_clean_limit_reached(): + log.info("岗位清洗:已达累计上限 {} 条,停止", settings.clean_total_limit) + _stop_event.set() + return + + +async def _fetch_one() -> dict | None: + """从 PG 锁定一条待清洗数据并标记为 cleaning""" async with PgSession() as pg: result = await pg.execute( - text(""" - SELECT * FROM app_job_data - WHERE clean_status = 'pending' AND recruit_category in (0,1,2) - LIMIT :limit - FOR UPDATE SKIP LOCKED - """), - {"limit": settings.clean_batch_size}, - ) - rows = result.mappings().all() - if not rows: - return - - ids = [r["id"] for r in rows] - await pg.execute( text(""" UPDATE app_job_data SET clean_status = 'cleaning', clean_started_at = NOW() - WHERE id = ANY(:ids) + WHERE id = ( + SELECT id FROM app_job_data + WHERE clean_status = 'pending' AND recruit_category IN (0,1,2) + LIMIT 1 + FOR UPDATE SKIP LOCKED + ) + RETURNING * """), - {"ids": ids}, ) - await pg.commit() - - log.info("岗位清洗:锁定{}条数据", len(rows)) - - # 2. 协程并发清洗,信号量限流 - sem = asyncio.Semaphore(settings.clean_concurrency) - tasks = [_clean_one(sem, dict(r)) for r in rows] - results = await asyncio.gather(*tasks, return_exceptions=True) - - # 汇总 - errors = sum(1 for r in results if isinstance(r, Exception)) - _cleaned_total += len(rows) - log.info("岗位清洗:本批完成,共{}条,异常{}条,累计{}条", len(rows), errors, _cleaned_total) - - if is_clean_limit_reached(): - log.info("岗位清洗:已达累计上限 {} 条,停止清洗任务", settings.clean_total_limit) - - -async def _clean_one(sem: asyncio.Semaphore, data: dict) -> None: - """单条岗位清洗""" - async with sem: - try: - await _do_clean(data) - except Exception as e: - log.error("岗位清洗异常, id={}: {}", data["id"], e) - # 保持 cleaning 状态,由僵尸恢复任务重置 + row = result.mappings().first() + if row: + await pg.commit() + return dict(row) + return None async def _do_clean(data: dict) -> None: @@ -253,7 +270,6 @@ async def _extract_skill_tags(job_id: int, result: dict) -> None: if not name or len(name) > 50: continue - # 每个 tag 单独 session,避免死锁 real_id = await _find_or_create_skill_tag(name) if real_id and real_id not in tag_ids: tag_ids.append(real_id) @@ -311,7 +327,7 @@ async def _find_or_create_company(short_name: str, urllistid: int | None = None) ) await mysql.commit() - # 锁外处理 logo:仅新建公司时执行,上传是网络IO,不阻塞其他协程;失败不影响主流程 + # 锁外处理 logo try: logo_b64 = await _get_logo_base64(urllistid) if logo_b64: