From fda30cd2959701b3e0b9acbd0920134d84a1f0f5 Mon Sep 17 00:00:00 2001 From: zk Date: Tue, 23 Jun 2026 10:03:51 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E6=9C=80=E5=A4=A7=E6=95=B0?= =?UTF-8?q?=E9=87=8F=E9=99=90=E5=88=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env | 13 +++++++------ app/config/settings.py | 1 + app/services/job_clean_service.py | 20 +++++++++++++++++++- 3 files changed, 27 insertions(+), 7 deletions(-) diff --git a/.env b/.env index c5a9f1e..913a6a5 100644 --- a/.env +++ b/.env @@ -8,10 +8,10 @@ PG_PASSWORD=feAeyR0u2fJGSS5ooFdHnSbyHQNY4WlV PG_DB=postgres # MySQL(业务库) -DB_HOST=192.168.31.105 -DB_PORT=3306 +DB_HOST=8.163.14.142 +DB_PORT=30006 DB_USER=root -DB_PASSWORD=123456 +DB_PASSWORD=^CgDatabase2020 DB_NAME=offerpie @@ -24,9 +24,10 @@ ANTHROPIC_API_KEY=sk-43ccdb29caa7e9ebe0db8ac0958c63f6d3a2d62e59064d3d26d94332055 ANTHROPIC_BASE_URL=https://code.warpdevloper.cloud # 岗位清洗参数 -CLEAN_BATCH_SIZE=10 -CLEAN_CONCURRENCY=10 -CLEAN_INTERVAL_SECONDS=180 +CLEAN_BATCH_SIZE=20 +CLEAN_CONCURRENCY=20 +CLEAN_INTERVAL_SECONDS=100 +CLEAN_TOTAL_LIMIT=0 # 公司补充参数 COMPANY_BATCH_SIZE=20 diff --git a/app/config/settings.py b/app/config/settings.py index 321c82b..e95b101 100644 --- a/app/config/settings.py +++ b/app/config/settings.py @@ -38,6 +38,7 @@ class Settings(BaseSettings): clean_batch_size: int = 100 clean_concurrency: int = 80 clean_interval_seconds: int = 200 + clean_total_limit: int = 0 # 累计清洗总数上限,达到后停止清洗任务;0 = 不限制 # ──────────── 公司补充参数 ──────────── company_batch_size: int = 20 diff --git a/app/services/job_clean_service.py b/app/services/job_clean_service.py index 57575bb..f7fb0a4 100644 --- a/app/services/job_clean_service.py +++ b/app/services/job_clean_service.py @@ -24,9 +24,23 @@ _id_gen = SnowflakeGenerator(instance=1) # 公司创建锁(防止并发重复插入同一公司) _company_lock = asyncio.Lock() +# 累计已清洗总数(用于总数上限控制,进程内统计) +_cleaned_total = 0 + + +def is_clean_limit_reached() -> bool: + """是否已达到累计清洗总数上限(clean_total_limit=0 表示不限制)""" + return 0 < settings.clean_total_limit <= _cleaned_total + async def run_job_clean() -> None: """一次批量清洗任务""" + global _cleaned_total + + # 总数上限:已达上限直接跳过 + if is_clean_limit_reached(): + return + # 1. 从 PG 锁定一批待清洗数据 async with PgSession() as pg: result = await pg.execute( @@ -62,7 +76,11 @@ async def run_job_clean() -> None: # 汇总 errors = sum(1 for r in results if isinstance(r, Exception)) - log.info("岗位清洗:本批完成,共{}条,异常{}条", len(rows), errors) + _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: