Compare commits

...
15 Commits
Author SHA1 Message Date
zk bdbbf32b8f 修改模型 2026-08-05 17:30:35 +08:00
zk c5440eefa9 修改模型 2026-07-30 18:00:09 +08:00
zk ca227ebef7 配置连接池 2026-07-15 16:14:31 +08:00
zk 54e415b5aa 冗余公司名字 2026-07-15 12:20:06 +08:00
zk 3459090284 扩大技能提示词,重构提取逻辑,岗位数据小于简历端数据,思路是简历覆盖岗位技能表 2026-07-14 18:29:27 +08:00
zk 0e10432583 缓存过滤省份 2026-07-13 20:07:00 +08:00
zk 0bb5791cd5 添加过滤条件 2026-07-13 18:43:07 +08:00
zk dfb0f177e2 放弃探针方案,但协程独立重试 2026-07-09 19:21:35 +08:00
zk 64bf6d7af0 修改探针逻辑 2026-07-09 15:47:19 +08:00
zk 900cf57bb3 方案修改为 多线程 滚动 2026-07-09 15:30:02 +08:00
zk f50a28bb35 固定间隔时间 2026-07-09 10:58:53 +08:00
zk 1a270d2c67 修改环境配置 2026-07-08 10:57:45 +08:00
zk ca0b13a45d 添加ai日志 2026-07-07 21:43:49 +08:00
zk cb94fb7bfe 添加AI门闸 2026-07-07 21:09:23 +08:00
zk 5cebc1f929 补充公司log 2026-07-03 17:27:18 +08:00
18 changed files with 509 additions and 191 deletions
+26 -6
View File
@@ -6,6 +6,8 @@ PG_PORT=5432
PG_USER=postgres
PG_PASSWORD=feAeyR0u2fJGSS5ooFdHnSbyHQNY4WlV
PG_DB=postgres
PG_POOL_SIZE=10
PG_MAX_OVERFLOW=10
# MySQL(业务库)
DB_HOST=8.163.14.142
@@ -13,6 +15,11 @@ DB_PORT=30006
DB_USER=root
DB_PASSWORD=^CgDatabase2020
DB_NAME=offerpie
MYSQL_POOL_SIZE=10
MYSQL_MAX_OVERFLOW=10
# 连接池获取连接等待超时(秒),两库共用
DB_POOL_TIMEOUT=30
# AI 供应商
@@ -23,19 +30,32 @@ VOLCENGINE_BASE_URL=https://ark.cn-beijing.volces.com/api/v3
ANTHROPIC_API_KEY=sk-43ccdb29caa7e9ebe0db8ac0958c63f6d3a2d62e59064d3d26d94332055a9bc9
ANTHROPIC_BASE_URL=https://code.warpdevloper.cloud
# Claude 备用中转(Anthropic 风格)
ANTHROPIC_ALT_API_KEY=sk-ac5544257dcf843de887fa153c4f4a11fa5007cbae4b8f4eec895a354c89ff7b
ANTHROPIC_ALT_BASE_URL=http://98.142.250.143:8080
# OpenAIOpenAI 兼容风格)
OPENAI_API_KEY=sk-12663c3a1ae5db0bc9aee78ffa64e7a0b95bcf7a15541040dd2bce573130f5c8
OPENAI_BASE_URL=https://re.94xy.cn/v1
# 岗位清洗参数
CLEAN_BATCH_SIZE=20
CLEAN_CONCURRENCY=20
CLEAN_INTERVAL_SECONDS=100
CLEAN_WORKER_COUNT=20
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
# 日志
LOGGING_LEVEL=INFO
# 阿里云 OSS
OSS_ACCESS_KEY_ID=LTAI5tEdLKKQUKhTyUpfH5Mk
OSS_ACCESS_KEY_SECRET=RjUdTrq0V5qA4b3BUElNhXqs3ZLp5k
OSS_ENDPOINT=oss-cn-guangzhou.aliyuncs.com
OSS_BUCKET=offerpie
OSS_DOMAIN=https://offerpie.oss-cn-guangzhou.aliyuncs.com
+32 -6
View File
@@ -6,6 +6,8 @@ PG_PORT=5432
PG_USER=postgres
PG_PASSWORD=feAeyR0u2fJGSS5ooFdHnSbyHQNY4WlV
PG_DB=postgres
PG_POOL_SIZE=10
PG_MAX_OVERFLOW=10
# MySQL(业务库)
DB_HOST=8.163.14.142
@@ -13,23 +15,47 @@ DB_PORT=30006
DB_USER=root
DB_PASSWORD=^CgDatabase2020
DB_NAME=offerpie
MYSQL_POOL_SIZE=10
MYSQL_MAX_OVERFLOW=10
# 连接池获取连接等待超时(秒),两库共用
DB_POOL_TIMEOUT=30
# AI 供应商
VOLCENGINE_API_KEY=fd065993-bee2-4f31-8bf2-56d5d3012c02
VOLCENGINE_BASE_URL=https://ark.cn-beijing.volces.com/api/v3
# ClaudeAnthropic 风格)
ANTHROPIC_API_KEY=sk-43ccdb29caa7e9ebe0db8ac0958c63f6d3a2d62e59064d3d26d94332055a9bc9
ANTHROPIC_BASE_URL=https://code.warpdevloper.cloud
# Claude 备用中转(Anthropic 风格)
ANTHROPIC_ALT_API_KEY=sk-ac5544257dcf843de887fa153c4f4a11fa5007cbae4b8f4eec895a354c89ff7b
ANTHROPIC_ALT_BASE_URL=http://98.142.250.143:8080
# OpenAIOpenAI 兼容风格)
OPENAI_API_KEY=sk-12663c3a1ae5db0bc9aee78ffa64e7a0b95bcf7a15541040dd2bce573130f5c8
OPENAI_BASE_URL=https://re.94xy.cn/v1
# 岗位清洗参数
CLEAN_BATCH_SIZE=100
CLEAN_CONCURRENCY=80
CLEAN_INTERVAL_SECONDS=200
CLEAN_WORKER_COUNT=20
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
# 日志
LOGGING_LEVEL=INFO
# 阿里云 OSS
OSS_ACCESS_KEY_ID=LTAI5tEdLKKQUKhTyUpfH5Mk
OSS_ACCESS_KEY_SECRET=RjUdTrq0V5qA4b3BUElNhXqs3ZLp5k
OSS_ENDPOINT=oss-cn-guangzhou.aliyuncs.com
OSS_BUCKET=offerpie
OSS_DOMAIN=https://offerpie.oss-cn-guangzhou.aliyuncs.com
+26 -6
View File
@@ -6,6 +6,8 @@ PG_PORT=5432
PG_USER=postgres
PG_PASSWORD=feAeyR0u2fJGSS5ooFdHnSbyHQNY4WlV
PG_DB=postgres
PG_POOL_SIZE=10
PG_MAX_OVERFLOW=10
# MySQL(业务库)
DB_HOST=192.168.31.105
@@ -13,24 +15,42 @@ DB_PORT=3306
DB_USER=root
DB_PASSWORD=123456
DB_NAME=offerpie
MYSQL_POOL_SIZE=10
MYSQL_MAX_OVERFLOW=10
# 连接池获取连接等待超时(秒),两库共用
DB_POOL_TIMEOUT=30
# AI 供应商
VOLCENGINE_API_KEY=fd065993-bee2-4f31-8bf2-56d5d3012c02
VOLCENGINE_BASE_URL=https://ark.cn-beijing.volces.com/api/v3
# Claude 备用中转(Anthropic 风格)
ANTHROPIC_ALT_API_KEY=sk-ac5544257dcf843de887fa153c4f4a11fa5007cbae4b8f4eec895a354c89ff7b
ANTHROPIC_ALT_BASE_URL=http://98.142.250.143:8080
# OpenAIOpenAI 兼容风格)
OPENAI_API_KEY=sk-12663c3a1ae5db0bc9aee78ffa64e7a0b95bcf7a15541040dd2bce573130f5c8
OPENAI_BASE_URL=https://re.94xy.cn/v1
# 岗位清洗参数
CLEAN_BATCH_SIZE=100
CLEAN_CONCURRENCY=50
CLEAN_INTERVAL_SECONDS=180
CLEAN_WORKER_COUNT=20
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
# 日志
LOGGING_LEVEL=INFO
# 阿里云 OSS
OSS_ACCESS_KEY_ID=LTAI5tEdLKKQUKhTyUpfH5Mk
OSS_ACCESS_KEY_SECRET=RjUdTrq0V5qA4b3BUElNhXqs3ZLp5k
OSS_ENDPOINT=oss-cn-guangzhou.aliyuncs.com
OSS_BUCKET=offerpie
OSS_DOMAIN=https://offerpie.oss-cn-guangzhou.aliyuncs.com
+4 -4
View File
@@ -9,14 +9,14 @@ from app.ai.models import LLM
class JobCleanModel:
"""岗位清洗模块"""
# 第一次AI:结构化提取岗位信息
STRUCTURE = LLM.CLAUDE_OPUS.create(temperature=0)
STRUCTURE = LLM.CLAUDE_SONNET_4_5.create(temperature=0)
# 第二次AI:专业匹配
MAJOR_MATCH = LLM.CLAUDE_OPUS.create(temperature=0)
MAJOR_MATCH = LLM.CLAUDE_SONNET_4_5.create(temperature=0)
# 第三次AI:技能提取
SKILL_EXTRACT = LLM.CLAUDE_OPUS.create(temperature=0)
SKILL_EXTRACT = LLM.CLAUDE_SONNET_4_5.create(temperature=0)
class CompanyCleanModel:
"""公司补充模块"""
# 公司信息补充
ENRICH = LLM.CLAUDE_OPUS.create(temperature=0)
ENRICH = LLM.CLAUDE_SONNET_4_5.create(temperature=0)
+6
View File
@@ -18,6 +18,8 @@ from app.config import settings
# 供应商连接配置 = (api_key函数, base_url函数)
_VOLCENGINE = (lambda: settings.volcengine_api_key, lambda: settings.volcengine_base_url)
_ANTHROPIC = (lambda: settings.anthropic_api_key, lambda: settings.anthropic_base_url)
_ANTHROPIC_ALT = (lambda: settings.anthropic_alt_api_key, lambda: settings.anthropic_alt_base_url)
_OPENAI = (lambda: settings.openai_api_key, lambda: settings.openai_base_url)
class LLM(Enum):
@@ -32,6 +34,10 @@ class LLM(Enum):
# ClaudeAnthropic 风格)
CLAUDE_OPUS = ("claude-opus-4-6", ChatAnthropic, *_ANTHROPIC)
CLAUDE_SONNET_4_5 = ("claude-sonnet-4-5-20250929", ChatAnthropic, *_ANTHROPIC_ALT)
# OpenAIOpenAI 兼容)
GPT_5_4_MINI = ("gpt-5.4-mini", ChatOpenAI, *_OPENAI)
def __init__(self, model_name: str, cls, api_key_fn, base_url_fn):
self.model_name = model_name
+25 -9
View File
@@ -50,18 +50,34 @@ MAJOR_MATCH_SYSTEM = """你是一个岗位专业匹配助手。根据岗位信
# ──────────── 第三次AI:技能提取 ────────────
SKILL_EXTRACT_SYSTEM = """你是一个技能提取助手。根据岗位信息,提取该岗位要求的核心专业能力和工具技能。
# 说明:与 Java 简历侧(UserProfileAnalyzeService.extractSkillTags)配对使用。
# 策略「岗位精、简历广」:岗位侧只提该岗位最核心、最不可或缺的技能(≤10),简历侧尽可能广地扩散提取。
# 只要候选人的宽技能集覆盖了岗位的核心要求即视为匹配。评估看「岗位核心被简历覆盖率」而非交集。
SKILL_EXTRACT_SYSTEM = """你是一个技能提取助手。根据岗位信息,只提取该岗位最核心、最不可或缺的硬技能(无论岗位是否明确写出,都抓最本质的那几个)。
返回JSON数组格式,如:["java", "spring boot", "mysql", "redis"]
规则:
1. 统一使用小写字母
1. 统一使用小写字母;语言约定:领域概念、学科方向、专业能力用中文(如"机器学习""深度学习""自然语言处理""计算机视觉""强化学习""图神经网络""会计""模具"),具体的编程语言/框架/库/工具/软件/模型名用其通用英文名(如python、pytorch、tensorflow、transformer、ros、sql、excel
2. 只保留核心词,去掉多余修饰(如"plc编程""plc""c语言""c""cad制图""cad"
3. 同一技能只保留最具体的表述,不要同时出现上位词和下位词(如有"机械设计"就不要再出"机械"
4. 提取范围包括:技术栈、专业领域知识、行业工具、专业资质能力等
5. 提取纯软技能(如沟通能力、团队协作、学习能力、积极主动
6. 不提取过于宽泛的标签(如"办公软件""windows"
7. 如果岗位完全没有专业能力要求(纯看态度和素质),返回空数组 []
8. 最多15个,按重要性排序
9. 只返回JSON数组,不要其他内容"""
3. 归一到常用统一粒度:同一技能只保留一个最通用常见的表述(如"模具设计"统一为"模具""前端开发"统一为"前端",不要同时出现"机械设计""机械"
4. 优先使用具体工具/技术名,不要用抽象概念词(如用"kafka"不用"消息队列",用"redis"不用"缓存",用"mysql"不用"数据库"
5. 提取具体、有明确名称的硬技能:编程语言、框架、软件、工具、平台、专业资质证书、明确的技术/工艺名称,以及通用学科名(如"会计""模具"
6. 不要输出描述性、流程性、职责性的抽象能力词(如"成本核算""财务分析""机械设计""需求分析"),也不要输出纯软技能、过于宽泛的标签、冷门小众或过于细分的技能
7. 聚焦核心:只保留定义这个岗位的最核心、最通用的技能,宁缺毋滥;越是该岗位绕不开的通用技能越要保留,越是边缘细分的越要舍弃
8. 最多提取10个,按重要性排序
9. 如果岗位完全没有专业能力要求(纯看态度和素质),返回空数组 []
10. 只返回JSON数组,不要其他内容
示例1(技术岗,只留核心):
输入:负责后端接口开发,要求熟悉Spring Boot、MyBatis、Redis
输出:["java", "spring boot", "mysql", "redis", "sql"]
示例2(财务岗,只留核心具体硬技能):
输入:负责费用管理与审核,会计学相关专业,熟练使用SAP
输出:["会计", "sap", "excel"]
示例3(制造岗,只留核心具体硬技能):
输入:负责模具开发,要求熟悉SolidWorks三维建模
输出:["模具", "solidworks", "cad"]
示例4(无专业能力):
输入:要求具备较强的沟通能力和创新意识,积极主动
输出:[]"""
# ──────────── 公司补充 ────────────
+27 -8
View File
@@ -14,7 +14,7 @@ class Settings(BaseSettings):
pg_password: str = ""
pg_db: str = "postgres"
pg_pool_size: int = 10
pg_max_overflow: int = 20
pg_max_overflow: int = 10
# ──────────── MySQL(业务库)────────────
db_host: str = "192.168.31.105"
@@ -23,7 +23,10 @@ class Settings(BaseSettings):
db_password: str = "123456"
db_name: str = "offerpie"
mysql_pool_size: int = 10
mysql_max_overflow: int = 20
mysql_max_overflow: int = 10
# 连接池获取连接的等待超时(秒),两个数据源共用
db_pool_timeout: int = 30
# ──────────── AI 供应商 ────────────
# 火山引擎(OpenAI 兼容风格)
@@ -34,20 +37,36 @@ class Settings(BaseSettings):
anthropic_api_key: str = "sk-43ccdb29caa7e9ebe0db8ac0958c63f6d3a2d62e59064d3d26d94332055a9bc9"
anthropic_base_url: str = "https://code.warpdevloper.cloud"
# Claude 备用中转(Anthropic 风格)
anthropic_alt_api_key: str = "sk-ac5544257dcf843de887fa153c4f4a11fa5007cbae4b8f4eec895a354c89ff7b"
anthropic_alt_base_url: str = "http://98.142.250.143:8080"
# OpenAIOpenAI 兼容风格)
openai_api_key: str = "sk-12663c3a1ae5db0bc9aee78ffa64e7a0b95bcf7a15541040dd2bce573130f5c8"
openai_base_url: str = "https://re.94xy.cn/v1"
# ──────────── 岗位清洗参数 ────────────
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
# ──────────── 阿里云 OSS ────────────
oss_access_key_id: str = ""
oss_access_key_secret: str = ""
oss_endpoint: str = "oss-cn-guangzhou.aliyuncs.com"
oss_bucket: str = "offerpie"
# 固定上传目录(对象 key 前缀)
oss_upload_dir: str = "company/logo"
# 下载访问域名
oss_domain: str = "https://offerpie.oss-cn-guangzhou.aliyuncs.com"
# ──────────── 日志 ────────────
logging_level: str = "INFO"
log_file_name: str = "cleaner.log"
+4
View File
@@ -39,7 +39,9 @@ async def init_db() -> None:
settings.pg_url,
pool_size=settings.pg_pool_size,
max_overflow=settings.pg_max_overflow,
pool_timeout=settings.db_pool_timeout,
pool_recycle=3600,
pool_pre_ping=True,
echo=False,
)
_pg_session_factory = async_sessionmaker(_pg_engine, expire_on_commit=False)
@@ -48,7 +50,9 @@ async def init_db() -> None:
settings.mysql_url,
pool_size=settings.mysql_pool_size,
max_overflow=settings.mysql_max_overflow,
pool_timeout=settings.db_pool_timeout,
pool_recycle=3600,
pool_pre_ping=True,
echo=False,
)
_mysql_session_factory = async_sessionmaker(_mysql_engine, expire_on_commit=False)
+42
View File
@@ -0,0 +1,42 @@
"""阿里云 OSS 上传:往固定目录上传,返回下载地址"""
import asyncio
import base64
import os
from uuid import uuid4
import oss2
from app.config import settings
from app.core.logger import log
_auth = oss2.Auth(settings.oss_access_key_id, settings.oss_access_key_secret)
_bucket = oss2.Bucket(_auth, f"https://{settings.oss_endpoint}", settings.oss_bucket)
async def upload(data: bytes, file_name: str) -> str:
"""上传字节内容到固定目录,返回下载地址"""
key = f"{settings.oss_upload_dir}/{uuid4().hex[:18]}{os.path.splitext(file_name)[1]}"
await asyncio.to_thread(_bucket.put_object, key, data)
url = f"{settings.oss_domain}/{key}"
log.info("OSS 上传成功 [{}] -> {}", file_name, url)
return url
async def upload_base64(b64: str, file_name: str = "image.png") -> str | None:
"""上传 base64 图片到固定目录,返回下载地址;空值或解码失败返回 None
:param b64: base64 字符串,兼容 data URI 前缀(data:image/png;base64,xxx
:param file_name: 用于取后缀,默认 png
"""
if not b64 or not b64.strip():
return None
# 去掉 data URI 前缀
if "," in b64 and b64.strip().startswith("data:"):
b64 = b64.split(",", 1)[1]
try:
data = base64.b64decode(b64)
except Exception as e:
log.warning("base64 解码失败: {}", e)
return None
return await upload(data, file_name)
+20 -13
View File
@@ -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 + SIGTERMWindows: 仅靠 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__":
+22
View File
@@ -0,0 +1,22 @@
"""PostgreSQL: ai_call_log 表模型"""
from datetime import datetime
from typing import Optional
from sqlalchemy import BigInteger, DateTime, String, Text
from sqlalchemy.orm import Mapped, mapped_column
from app.core.database import PgBase
class AiCallLog(PgBase):
"""AI 调用日志表"""
__tablename__ = "ai_call_log"
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, comment="主键ID(雪花)")
scene: Mapped[str] = mapped_column(String(32), nullable=False, comment="场景标识: structure/major_match/skill_extract/company_enrich")
system_prompt: Mapped[str] = mapped_column(Text, nullable=False, comment="系统提示词")
user_message: Mapped[str] = mapped_column(Text, nullable=False, comment="发给AI的用户消息")
response: Mapped[Optional[str]] = mapped_column(Text, comment="AI返回原文")
created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, comment="创建时间")
+35
View File
@@ -0,0 +1,35 @@
"""PostgreSQL: app_url_list 表模型"""
from datetime import datetime
from typing import Optional
from sqlalchemy import BigInteger, DateTime, Integer, String, Text
from sqlalchemy.orm import Mapped, mapped_column
from app.core.database import PgBase
class AppUrlList(PgBase):
"""爬虫 URL 任务列表"""
__tablename__ = "app_url_list"
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True, comment="自增主键")
status: Mapped[str] = mapped_column(String(32), default="pending", nullable=False, comment="任务状态")
input_company_name: Mapped[str] = mapped_column(String(255), nullable=False, comment="输入公司名")
input_url: Mapped[Optional[str]] = mapped_column(String(1024), comment="输入URL")
output_url: Mapped[Optional[str]] = mapped_column(String(1024), comment="输出URL")
post_login: Mapped[Optional[int]] = mapped_column(Integer, comment="发布登录")
crawl_login: Mapped[Optional[int]] = mapped_column(Integer, comment="爬取登录")
type_sz: Mapped[Optional[int]] = mapped_column(Integer, comment="类型sz")
type_sx: Mapped[Optional[int]] = mapped_column(Integer, comment="类型sx")
crawlok_at: Mapped[Optional[str]] = mapped_column(String(32), comment="-2没有数据或者数据全部过期或者数据缺失 -3 过滤后没数据, -4无法修复,")
pyname: Mapped[Optional[str]] = mapped_column(String(255), comment="脚本名字")
type_xz: Mapped[Optional[int]] = mapped_column(Integer, comment="类型xz")
error_message: Mapped[Optional[str]] = mapped_column(Text, comment="错误信息")
logo: Mapped[Optional[str]] = mapped_column(Text, comment="公司logo (base64编码的64x64 PNG图片)")
started_at: Mapped[Optional[datetime]] = mapped_column(DateTime, comment="开始时间")
finished_at: Mapped[Optional[datetime]] = mapped_column(DateTime, comment="完成时间")
created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, comment="创建时间")
updated_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, comment="更新时间")
output_company_name: Mapped[Optional[str]] = mapped_column(String(255), comment="输出公司名")
+8 -37
View File
@@ -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:
+57 -12
View File
@@ -1,14 +1,24 @@
"""AI 调用工具封装"""
"""AI 调用工具封装
每个协程独立重试,失败后 sleep 200ms 再试,直到成功。
不再有全局门闸,互不影响。
"""
import asyncio
import re
from typing import Any
from datetime import datetime
from typing import Any, Optional
from json_repair import repair_json
from langchain_core.language_models import BaseChatModel
from langchain_core.messages import SystemMessage, HumanMessage
from snowflake import SnowflakeGenerator
from app.core.logger import log
# AI 日志专用雪花ID
_log_id_gen = SnowflakeGenerator(instance=2)
# 匹配 <think>任意内容</think>,用于剥离推理模型的思考过程
_THINK_RE = re.compile(r"<think>.*?</think>", re.DOTALL | re.IGNORECASE)
@@ -28,19 +38,31 @@ def parse_llm_json(text: str) -> Any:
return repair_json(cleaned, return_objects=True)
async def ai_chat(llm: BaseChatModel, system_prompt: str, user_message: str) -> str:
"""异步调用 LLM,返回原始文本"""
messages = [
SystemMessage(content=system_prompt),
HumanMessage(content=user_message),
]
response = await llm.ainvoke(messages)
return response.content
async def ai_chat(llm: BaseChatModel, system_prompt: str, user_message: str, scene: Optional[str] = None) -> str:
"""异步调用 LLM,返回原始文本。失败后独立重试直到成功。"""
attempt = 0
while True:
attempt += 1
try:
messages = [
SystemMessage(content=system_prompt),
HumanMessage(content=user_message),
]
response = await llm.ainvoke(messages)
content = response.content
# 记录 AI 调用日志(失败不影响主流程)
if scene:
await _save_call_log(scene, system_prompt, user_message, content)
return content
except Exception as e:
log.warning("AI 调用失败(第{}次), 200ms 后重试: {}", attempt, e)
await asyncio.sleep(0.2)
async def ai_chat_json(llm: BaseChatModel, system_prompt: str, user_message: str) -> Any:
async def ai_chat_json(llm: BaseChatModel, system_prompt: str, user_message: str, scene: Optional[str] = None) -> Any:
"""异步调用 LLM,返回解析后的 JSON 对象"""
raw = await ai_chat(llm, system_prompt, user_message)
raw = await ai_chat(llm, system_prompt, user_message, scene=scene)
if not raw or not raw.strip():
log.warning("AI 返回为空")
return None
@@ -49,3 +71,26 @@ async def ai_chat_json(llm: BaseChatModel, system_prompt: str, user_message: str
except Exception as e:
log.warning("AI JSON 解析失败: {}, raw={}", e, raw[:200])
return None
async def _save_call_log(scene: str, system_prompt: str, user_message: str, response: Optional[str]) -> None:
"""异步写入 AI 调用日志到 PG,失败只 warning 不影响主流程"""
try:
from app.core.database import PgSession
from sqlalchemy import insert
from app.models.pg.ai_call_log import AiCallLog
async with PgSession() as pg:
await pg.execute(
insert(AiCallLog).values(
id=next(_log_id_gen),
scene=scene,
system_prompt=system_prompt,
user_message=user_message,
response=response,
created_at=datetime.now(),
)
)
await pg.commit()
except Exception as e:
log.warning("AI 调用日志写入失败: {}", e)
+51 -29
View File
@@ -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:
@@ -62,7 +84,7 @@ async def _do_clean(company: dict) -> None:
short_name = company.get("short_name", "")
user_msg = f"【公司简称】\n{short_name}\n\n【行业列表】\n{dict_cache.industry_text}"
result = await ai_chat_json(CompanyCleanModel.ENRICH, COMPANY_ENRICH_SYSTEM, user_msg)
result = await ai_chat_json(CompanyCleanModel.ENRICH, COMPANY_ENRICH_SYSTEM, user_msg, scene="company_enrich")
if result is None or not result.get("valid", False):
await _update_status(company_id, 4)
+3 -2
View File
@@ -64,9 +64,10 @@ class DictCacheService:
for m in majors
)
# 地区(市级)
# 地区(市级province_code 不为空且 city_code 为空;直辖市取市级避免与省级同名重复
result = await session.execute(text("""
SELECT code, name FROM bg_china_regions_code WHERE city_code IS NULL
SELECT code, name FROM bg_china_regions_code
WHERE city_code IS NULL AND province_code IS NOT NULL
"""))
self._region_list = [dict(r) for r in result.mappings().all()]
+119 -58
View File
@@ -1,4 +1,8 @@
"""岗位清洗服务(协程版)"""
"""岗位清洗服务(持续消费模型)
启动 N 个 worker 协程,每个 worker 循环:取一条 → 处理 → 取下一条。
没数据时短暂休眠后重试。
"""
import asyncio
import json
@@ -10,6 +14,7 @@ from sqlalchemy import text, insert
from app.config import settings
from app.core.database import PgSession, MysqlSession
from app.core.logger import log
from app.core.oss import upload_base64
from app.ai.model_config import JobCleanModel
from app.ai.prompts import JOB_STRUCTURE_SYSTEM, MAJOR_MATCH_SYSTEM, SKILL_EXTRACT_SYSTEM
from app.models.mysql.job import Job
@@ -27,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 verify_status = 'success' 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:
@@ -125,8 +143,8 @@ async def _do_clean(data: dict) -> None:
return
# 第一次AI:结构化提取
user_message = _build_user_message(data)
result = await ai_chat_json(JobCleanModel.STRUCTURE, JOB_STRUCTURE_SYSTEM, user_message)
user_message = await _build_user_message(data)
result = await ai_chat_json(JobCleanModel.STRUCTURE, JOB_STRUCTURE_SYSTEM, user_message, scene="structure")
if result is None or not result.get("valid", False):
log.info("[id={}] 丢弃:AI判定无效", data_id)
await _update_pg_status(data_id, "discarded")
@@ -134,7 +152,7 @@ async def _do_clean(data: dict) -> None:
# 公司处理
company_short_name = result.get("companyShortName") or data.get("company") or ""
company_id = await _find_or_create_company(company_short_name)
company_id = await _find_or_create_company(company_short_name, data.get("urllistid"))
# 地区处理
region_codes = []
@@ -216,7 +234,7 @@ async def _match_major(job_id: int, result: dict) -> None:
req = result.get("requirement", "")
user_msg = f"【岗位信息】\n标题: {title}\n职责: {desc}\n要求: {req}\n\n【专业分类列表】\n{dict_cache.major_category_text}"
data = await ai_chat_json(JobCleanModel.MAJOR_MATCH, MAJOR_MATCH_SYSTEM, user_msg)
data = await ai_chat_json(JobCleanModel.MAJOR_MATCH, MAJOR_MATCH_SYSTEM, user_msg, scene="major_match")
if data is None:
return
@@ -241,7 +259,7 @@ async def _extract_skill_tags(job_id: int, result: dict) -> None:
req = result.get("requirement", "")
user_msg = f"【岗位信息】\n标题: {title}\n职责: {desc}\n要求: {req}"
skills = await ai_chat_json(JobCleanModel.SKILL_EXTRACT, SKILL_EXTRACT_SYSTEM, user_msg)
skills = await ai_chat_json(JobCleanModel.SKILL_EXTRACT, SKILL_EXTRACT_SYSTEM, user_msg, scene="skill_extract")
if not skills or not isinstance(skills, list):
return
@@ -252,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)
@@ -284,8 +301,8 @@ async def _find_or_create_skill_tag(name: str) -> int | None:
return row.scalar()
async def _find_or_create_company(short_name: str) -> int:
"""查找或创建公司(加锁防并发重复)"""
async def _find_or_create_company(short_name: str, urllistid: int | None = None) -> int:
"""查找或创建公司(加锁防并发重复);新建公司时若有 logo 则上传 OSS 并回填地址"""
async with _company_lock:
async with MysqlSession() as mysql:
row = await mysql.execute(
@@ -309,7 +326,35 @@ async def _find_or_create_company(short_name: str) -> int:
)
)
await mysql.commit()
return company_id
# 锁外处理 logo
try:
logo_b64 = await _get_logo_base64(urllistid)
if logo_b64:
logo_url = await upload_base64(logo_b64, "logo.png")
if logo_url:
async with MysqlSession() as mysql:
await mysql.execute(
text("UPDATE bg_company SET logo_url = :url, update_time = :t WHERE id = :id"),
{"url": logo_url, "t": datetime.now(), "id": company_id},
)
await mysql.commit()
except Exception as e:
log.warning("[company={}] logo 上传失败: {}", company_id, e)
return company_id
async def _get_logo_base64(urllistid: int | None) -> str | None:
"""从 PG app_url_list 按 urllistid 读取 logobase64"""
if not urllistid:
return None
async with PgSession() as pg:
row = await pg.execute(
text("SELECT logo FROM app_url_list WHERE id = :id"),
{"id": urllistid},
)
return row.scalar()
async def _update_pg_status(data_id: int, status: str) -> None:
@@ -328,14 +373,18 @@ async def _update_pg_status(data_id: int, status: str) -> None:
await pg.commit()
def _build_user_message(data: dict) -> str:
async def _build_user_message(data: dict) -> str:
"""构建第一次AI的用户消息"""
# 从 app_url_list 取来源公司名,作为 companyShortName 的参考
input_company_name = await _get_input_company_name(data.get("urllistid"))
parts = [
"【原始数据】",
f"岗位名称: {data.get('job_title') or ''}",
f"薪资: {data.get('salary') or ''}",
f"工作地点: {data.get('location') or ''}",
f"公司: {data.get('company') or ''}",
f"来源公司名(参考): {input_company_name or ''}",
f"经验要求: {data.get('experience') or ''}",
f"学历要求: {data.get('education') or ''}",
f"岗位详情: {data.get('description') or ''}",
@@ -345,3 +394,15 @@ def _build_user_message(data: dict) -> str:
f"【行业列表】\n{dict_cache.industry_text}",
]
return "\n".join(parts)
async def _get_input_company_name(urllistid: int | None) -> str | None:
"""从 PG app_url_list 按 urllistid 读取 input_company_name"""
if not urllistid:
return None
async with PgSession() as pg:
row = await pg.execute(
text("SELECT input_company_name FROM app_url_list WHERE id = :id"),
{"id": urllistid},
)
return row.scalar()
+2 -1
View File
@@ -13,4 +13,5 @@ langchain-core>=0.3
# 工具
loguru>=0.7
snowflake-id>=1.0
json-repair>=0.30
json-repair>=0.30
oss2>=2.18