封装清洗过程

This commit is contained in:
zk
2026-07-28 15:40:23 +08:00
parent 20efceb5df
commit 6516944b69
5 changed files with 245 additions and 8 deletions
+3 -3
View File
@@ -1,10 +1,10 @@
ENV=dev
# MySQL(业务库)
DB_HOST=8.163.14.142
DB_PORT=30006
DB_HOST=192.168.31.105
DB_PORT=3306
DB_USER=root
DB_PASSWORD=^CgDatabase2020
DB_PASSWORD=123456
DB_NAME=offerpie
MYSQL_POOL_SIZE=10
MYSQL_MAX_OVERFLOW=10
+1 -1
View File
@@ -35,7 +35,7 @@ EXTRACT_SYSTEM_PROMPT = """你是招聘公告信息提取专家。用户会给
## 字段说明
- company_name:招聘的公司名称,取全称或公告中最完整的写法
- company_name:招聘的公司名称,取简称
- title:公告标题,如「XX公司2027届秋季校园招聘正式启动」
- company_intro:公司介绍段落,原文摘录,最多 300 字。公告中没有公司介绍时,可以根据你自己对该公司的了解补充一段简介;不了解该公司时填 null
- target_audience:面向对象,如「2027届本硕博」「2026届及2027届毕业生」,最多 200 字
+73
View File
@@ -0,0 +1,73 @@
"""公司业务服务:查找或创建公司基础记录,上传 Logo。"""
from __future__ import annotations
import threading
from datetime import datetime
from sqlalchemy import insert, text
from app.core.database import MysqlSession
from app.core.id_gen import next_id
from app.core.logger import log
from app.core.oss import upload_image_url
from app.models.company import Company
# 公司创建锁(防止并发重复插入同一公司)
_company_lock = threading.Lock()
def find_or_create_company(company_name: str, logo_url: str | None = None) -> int:
"""查找或创建公司,上传 Logo 并回填地址。
Args:
company_name: 公司名称(用于查重和创建)。
logo_url: Logo 图片原始地址,非空时下载并上传 OSS。
Returns:
公司 ID。
"""
if not company_name:
company_name = "未知公司"
with _company_lock:
with MysqlSession() as session:
row = session.execute(
text("SELECT id FROM bg_company WHERE short_name = :name LIMIT 1"),
{"name": company_name},
)
existing = row.scalar()
if existing:
return existing
company_id = next_id()
now = datetime.now()
session.execute(
insert(Company).values(
id=company_id,
name=company_name,
short_name=company_name,
status=0,
create_time=now,
update_time=now,
)
)
session.commit()
log.info("公司创建成功 [id={}]: {}", company_id, company_name)
# 锁外处理 Logo 上传
if logo_url:
try:
oss_url = upload_image_url(logo_url)
if oss_url:
with MysqlSession() as session:
session.execute(
text("UPDATE bg_company SET logo_url = :url, update_time = :t WHERE id = :id"),
{"url": oss_url, "t": datetime.now(), "id": company_id},
)
session.commit()
log.info("公司 Logo 上传成功 [id={}]: {}", company_id, oss_url)
except Exception as exc:
log.warning("公司 Logo 上传失败 [id={}]: {}", company_id, exc)
return company_id
+167 -4
View File
@@ -1,7 +1,170 @@
"""招聘公告业务服务。"""
"""招聘公告业务服务:编排公告爬取、AI 信息提取与数据库保存"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import insert, text
from app.ai.extract.announcement_extract import extract_announcement
from app.core.database import MysqlSession
from app.core.id_gen import next_id
from app.core.logger import log
from app.models.recruit_announcement import RecruitAnnouncement
from app.models.recruit_announcement_batch import RecruitAnnouncementBatch
from app.models.recruit_announcement_category import RecruitAnnouncementCategory
from app.models.recruit_announcement_city import RecruitAnnouncementCity
from app.models.recruit_announcement_education import RecruitAnnouncementEducation
from app.models.recruit_announcement_tag import RecruitAnnouncementTag
from app.models.recruit_announcement_year import RecruitAnnouncementYear
from app.service.company_service import find_or_create_company
from app.tool.page_extract import extract_page
class RecruitAnnouncementService:
"""编排公告爬取、AI 信息提取与数据库保存"""
def _parse_datetime(value: str | None) -> datetime | None:
"""将 yyyy-MM-dd HH:mm:ss 字符串解析为 datetime,失败返回 None"""
if not value:
return None
try:
return datetime.strptime(value, "%Y-%m-%d %H:%M:%S")
except (ValueError, TypeError):
return None
pass
def _save_announcement(announcement_id: int, company_id: int, url: str, data: dict) -> None:
"""保存公告主表和六张关联表,单事务提交。"""
now = datetime.now()
with MysqlSession() as session:
# 主表
session.execute(
insert(RecruitAnnouncement).values(
id=announcement_id,
company_id=company_id,
company_name=data.get("company_name") or "",
title=data.get("title") or "",
company_intro=data.get("company_intro"),
target_audience=data.get("target_audience"),
major_require=data.get("major_require"),
remark=data.get("remark"),
written_exam=data.get("written_exam"),
apply_start_time=_parse_datetime(data.get("apply_start_time")),
apply_end_time=_parse_datetime(data.get("apply_end_time")),
apply_end_desc=data.get("apply_end_desc"),
invite_code=data.get("invite_code"),
announcement_url=url,
apply_url=data.get("apply_url"),
apply_email=data.get("apply_email"),
source=data.get("source"),
publish_time=_parse_datetime(data.get("publish_time")),
clean_status=1,
status=1,
create_time=now,
update_time=now,
)
)
# 招聘届数
for year in data.get("recruit_years") or []:
session.execute(
insert(RecruitAnnouncementYear).values(
id=next_id(),
announcement_id=announcement_id,
recruit_year=int(year),
create_time=now,
)
)
# 批次
for batch in data.get("batches") or []:
session.execute(
insert(RecruitAnnouncementBatch).values(
id=next_id(),
announcement_id=announcement_id,
batch_name=str(batch)[:64],
create_time=now,
)
)
# 标签
for tag in data.get("tags") or []:
session.execute(
insert(RecruitAnnouncementTag).values(
id=next_id(),
announcement_id=announcement_id,
tag_name=str(tag)[:64],
create_time=now,
)
)
# 城市
for city in data.get("cities") or []:
session.execute(
insert(RecruitAnnouncementCity).values(
id=next_id(),
announcement_id=announcement_id,
city_name=str(city)[:64],
create_time=now,
)
)
# 岗位大类
for category in data.get("categories") or []:
session.execute(
insert(RecruitAnnouncementCategory).values(
id=next_id(),
announcement_id=announcement_id,
category_name=str(category)[:64],
create_time=now,
)
)
# 学历要求
for edu in data.get("educations") or []:
session.execute(
insert(RecruitAnnouncementEducation).values(
id=next_id(),
announcement_id=announcement_id,
education_name=str(edu)[:32],
create_time=now,
)
)
session.commit()
def process_announcement(url: str) -> None:
"""处理单条公告 URL 的完整流程。"""
# 1. URL 去重
with MysqlSession() as session:
row = session.execute(
text("SELECT id FROM bg_recruit_announcement WHERE announcement_url = :url LIMIT 1"),
{"url": url},
)
if row.scalar():
log.info("公告已存在,跳过: {}", url)
return
# 2. 页面内容提取
result = extract_page(url)
if not result.content or len(result.content) < 50:
log.info("公告页内容过短({}字),跳过: {}", len(result.content) if result.content else 0, url)
return
# 3. AI 信息提取
data = extract_announcement(result.content)
if data is None:
log.warning("AI 信息提取失败,跳过: {}", url)
return
# 4. 公司处理
company_name = data.get("company_name") or ""
company_id = find_or_create_company(company_name, result.logo_url)
# 5. 保存公告
announcement_id = next_id()
try:
_save_announcement(announcement_id, company_id, url, data)
log.info("公告入库成功 [id={}]: {} | 公司={}", announcement_id, data.get("title"), company_name)
except Exception as exc:
log.error("公告入库失败 [url={}]: {}", url, exc)
+1
View File
@@ -16,6 +16,7 @@ langchain-anthropic>=0.3
langchain-core>=0.3
# 工具
json-repair>=0.61
loguru>=0.7
oss2==2.19.1
snowflake-id>=1.0