Compare commits

..
2 Commits
Author SHA1 Message Date
zk fda30cd295 添加最大数量限制 2026-06-23 10:03:51 +08:00
zk 300b9d9dc5 ChatAnthropic 2026-06-22 20:21:41 +08:00
7 changed files with 68 additions and 26 deletions
+11 -6
View File
@@ -8,10 +8,10 @@ PG_PASSWORD=feAeyR0u2fJGSS5ooFdHnSbyHQNY4WlV
PG_DB=postgres PG_DB=postgres
# MySQL(业务库) # MySQL(业务库)
DB_HOST=192.168.31.105 DB_HOST=8.163.14.142
DB_PORT=3306 DB_PORT=30006
DB_USER=root DB_USER=root
DB_PASSWORD=123456 DB_PASSWORD=^CgDatabase2020
DB_NAME=offerpie DB_NAME=offerpie
@@ -19,10 +19,15 @@ DB_NAME=offerpie
VOLCENGINE_API_KEY=fd065993-bee2-4f31-8bf2-56d5d3012c02 VOLCENGINE_API_KEY=fd065993-bee2-4f31-8bf2-56d5d3012c02
VOLCENGINE_BASE_URL=https://ark.cn-beijing.volces.com/api/v3 VOLCENGINE_BASE_URL=https://ark.cn-beijing.volces.com/api/v3
# ClaudeAnthropic 风格)
ANTHROPIC_API_KEY=sk-43ccdb29caa7e9ebe0db8ac0958c63f6d3a2d62e59064d3d26d94332055a9bc9
ANTHROPIC_BASE_URL=https://code.warpdevloper.cloud
# 岗位清洗参数 # 岗位清洗参数
CLEAN_BATCH_SIZE=100 CLEAN_BATCH_SIZE=20
CLEAN_CONCURRENCY=50 CLEAN_CONCURRENCY=20
CLEAN_INTERVAL_SECONDS=180 CLEAN_INTERVAL_SECONDS=100
CLEAN_TOTAL_LIMIT=0
# 公司补充参数 # 公司补充参数
COMPANY_BATCH_SIZE=20 COMPANY_BATCH_SIZE=20
+4 -4
View File
@@ -9,14 +9,14 @@ from app.ai.models import LLM
class JobCleanModel: class JobCleanModel:
"""岗位清洗模块""" """岗位清洗模块"""
# 第一次AI:结构化提取岗位信息 # 第一次AI:结构化提取岗位信息
STRUCTURE = LLM.DOUBAO_SEED_LITE.create(temperature=0) STRUCTURE = LLM.CLAUDE_OPUS.create(temperature=0)
# 第二次AI:专业匹配 # 第二次AI:专业匹配
MAJOR_MATCH = LLM.DOUBAO_SEED_LITE.create(temperature=0) MAJOR_MATCH = LLM.CLAUDE_OPUS.create(temperature=0)
# 第三次AI:技能提取 # 第三次AI:技能提取
SKILL_EXTRACT = LLM.DOUBAO_SEED_LITE.create(temperature=0) SKILL_EXTRACT = LLM.CLAUDE_OPUS.create(temperature=0)
class CompanyCleanModel: class CompanyCleanModel:
"""公司补充模块""" """公司补充模块"""
# 公司信息补充 # 公司信息补充
ENRICH = LLM.DOUBAO_SEED_LITE.create(temperature=0) ENRICH = LLM.CLAUDE_OPUS.create(temperature=0)
+24 -12
View File
@@ -4,36 +4,48 @@ Usage:
from app.ai.models import LLM from app.ai.models import LLM
llm = LLM.DOUBAO_SEED_LITE.create(temperature=0) llm = LLM.DOUBAO_SEED_LITE.create(temperature=0)
llm = LLM.CLAUDE_OPUS.create(temperature=0)
""" """
from enum import Enum from enum import Enum
from langchain_anthropic import ChatAnthropic
from langchain_core.language_models import BaseChatModel
from langchain_openai import ChatOpenAI from langchain_openai import ChatOpenAI
from app.config import settings from app.config import settings
# 供应商连接配置 # 供应商连接配置 = (api_key函数, base_url函数)
_VOLCENGINE = (lambda: settings.volcengine_api_key, lambda: settings.volcengine_base_url) _VOLCENGINE = (lambda: settings.volcengine_api_key, lambda: settings.volcengine_base_url)
_ANTHROPIC = (lambda: settings.anthropic_api_key, lambda: settings.anthropic_base_url)
class LLM(Enum): class LLM(Enum):
"""所有可用模型,每个枚举值 = (模型名, api_key函数, base_url函数)""" """所有可用模型,每个枚举值 = (模型名, 封装类, api_key函数, base_url函数)"""
# 火山引擎 # 火山引擎OpenAI 兼容)
DOUBAO_PRO_32K = ("doubao-1-5-pro-32k-250115", *_VOLCENGINE) DOUBAO_PRO_32K = ("doubao-1-5-pro-32k-250115", ChatOpenAI, *_VOLCENGINE)
DOUBAO_LITE_32K = ("doubao-1-5-lite-32k-250115", *_VOLCENGINE) DOUBAO_LITE_32K = ("doubao-1-5-lite-32k-250115", ChatOpenAI, *_VOLCENGINE)
DOUBAO_SEED_LITE = ("doubao-seed-2-0-lite-260215", *_VOLCENGINE) DOUBAO_SEED_LITE = ("doubao-seed-2-0-lite-260215", ChatOpenAI, *_VOLCENGINE)
DOUBAO_SEED_PRO = ("doubao-seed-2-0-pro-260215", *_VOLCENGINE) DOUBAO_SEED_PRO = ("doubao-seed-2-0-pro-260215", ChatOpenAI, *_VOLCENGINE)
DEEPSEEK_V4_FLASH = ("deepseek-v4-flash-260425", *_VOLCENGINE) DEEPSEEK_V4_FLASH = ("deepseek-v4-flash-260425", ChatOpenAI, *_VOLCENGINE)
def __init__(self, model_name: str, api_key_fn, base_url_fn): # ClaudeAnthropic 风格)
CLAUDE_OPUS = ("claude-opus-4-6", ChatAnthropic, *_ANTHROPIC)
def __init__(self, model_name: str, cls, api_key_fn, base_url_fn):
self.model_name = model_name self.model_name = model_name
self._cls = cls
self._api_key_fn = api_key_fn self._api_key_fn = api_key_fn
self._base_url_fn = base_url_fn self._base_url_fn = base_url_fn
def create(self, **kwargs) -> ChatOpenAI: def create(self, **kwargs) -> BaseChatModel:
"""创建 LLM 实例,kwargs 透传给 ChatOpenAItemperature, max_tokens 等)""" """创建 LLM 实例,kwargs 透传给底层封装temperature, max_tokens 等)
return ChatOpenAI(
封装类(ChatOpenAI / ChatAnthropic)均实现 langchain BaseChatModel 接口,
对上层调用方完全透明。
"""
return self._cls(
model=self.model_name, model=self.model_name,
api_key=self._api_key_fn(), api_key=self._api_key_fn(),
base_url=self._base_url_fn(), base_url=self._base_url_fn(),
+6
View File
@@ -26,13 +26,19 @@ class Settings(BaseSettings):
mysql_max_overflow: int = 20 mysql_max_overflow: int = 20
# ──────────── AI 供应商 ──────────── # ──────────── AI 供应商 ────────────
# 火山引擎(OpenAI 兼容风格)
volcengine_api_key: str = "fd065993-bee2-4f31-8bf2-56d5d3012c02" volcengine_api_key: str = "fd065993-bee2-4f31-8bf2-56d5d3012c02"
volcengine_base_url: str = "https://ark.cn-beijing.volces.com/api/v3" volcengine_base_url: str = "https://ark.cn-beijing.volces.com/api/v3"
# ClaudeAnthropic 风格)
anthropic_api_key: str = "sk-43ccdb29caa7e9ebe0db8ac0958c63f6d3a2d62e59064d3d26d94332055a9bc9"
anthropic_base_url: str = "https://code.warpdevloper.cloud"
# ──────────── 岗位清洗参数 ──────────── # ──────────── 岗位清洗参数 ────────────
clean_batch_size: int = 100 clean_batch_size: int = 100
clean_concurrency: int = 80 clean_concurrency: int = 80
clean_interval_seconds: int = 200 clean_interval_seconds: int = 200
clean_total_limit: int = 0 # 累计清洗总数上限,达到后停止清洗任务;0 = 不限制
# ──────────── 公司补充参数 ──────────── # ──────────── 公司补充参数 ────────────
company_batch_size: int = 20 company_batch_size: int = 20
+3 -3
View File
@@ -4,7 +4,7 @@ import re
from typing import Any from typing import Any
from json_repair import repair_json from json_repair import repair_json
from langchain_openai import ChatOpenAI from langchain_core.language_models import BaseChatModel
from langchain_core.messages import SystemMessage, HumanMessage from langchain_core.messages import SystemMessage, HumanMessage
from app.core.logger import log from app.core.logger import log
@@ -28,7 +28,7 @@ def parse_llm_json(text: str) -> Any:
return repair_json(cleaned, return_objects=True) return repair_json(cleaned, return_objects=True)
async def ai_chat(llm: ChatOpenAI, system_prompt: str, user_message: str) -> str: async def ai_chat(llm: BaseChatModel, system_prompt: str, user_message: str) -> str:
"""异步调用 LLM,返回原始文本""" """异步调用 LLM,返回原始文本"""
messages = [ messages = [
SystemMessage(content=system_prompt), SystemMessage(content=system_prompt),
@@ -38,7 +38,7 @@ async def ai_chat(llm: ChatOpenAI, system_prompt: str, user_message: str) -> str
return response.content return response.content
async def ai_chat_json(llm: ChatOpenAI, system_prompt: str, user_message: str) -> Any: async def ai_chat_json(llm: BaseChatModel, system_prompt: str, user_message: str) -> Any:
"""异步调用 LLM,返回解析后的 JSON 对象""" """异步调用 LLM,返回解析后的 JSON 对象"""
raw = await ai_chat(llm, system_prompt, user_message) raw = await ai_chat(llm, system_prompt, user_message)
if not raw or not raw.strip(): if not raw or not raw.strip():
+19 -1
View File
@@ -24,9 +24,23 @@ _id_gen = SnowflakeGenerator(instance=1)
# 公司创建锁(防止并发重复插入同一公司) # 公司创建锁(防止并发重复插入同一公司)
_company_lock = asyncio.Lock() _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: async def run_job_clean() -> None:
"""一次批量清洗任务""" """一次批量清洗任务"""
global _cleaned_total
# 总数上限:已达上限直接跳过
if is_clean_limit_reached():
return
# 1. 从 PG 锁定一批待清洗数据 # 1. 从 PG 锁定一批待清洗数据
async with PgSession() as pg: async with PgSession() as pg:
result = await pg.execute( 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)) 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: async def _clean_one(sem: asyncio.Semaphore, data: dict) -> None:
+1
View File
@@ -7,6 +7,7 @@ apscheduler>=3.10
# AI # AI
langchain-openai>=0.3 langchain-openai>=0.3
langchain-anthropic>=0.3
langchain-core>=0.3 langchain-core>=0.3
# 工具 # 工具