diff --git a/.env b/.env index f121ae1..09c2ace 100644 --- a/.env +++ b/.env @@ -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 diff --git a/app/ai/extract/prompt.py b/app/ai/extract/prompt.py index eb876fc..d4f45cf 100644 --- a/app/ai/extract/prompt.py +++ b/app/ai/extract/prompt.py @@ -35,7 +35,7 @@ EXTRACT_SYSTEM_PROMPT = """你是招聘公告信息提取专家。用户会给 ## 字段说明 -- company_name:招聘的公司名称,取全称或公告中最完整的写法 +- company_name:招聘的公司名称,取简称 - title:公告标题,如「XX公司2027届秋季校园招聘正式启动」 - company_intro:公司介绍段落,原文摘录,最多 300 字。公告中没有公司介绍时,可以根据你自己对该公司的了解补充一段简介;不了解该公司时填 null - target_audience:面向对象,如「2027届本硕博」「2026届及2027届毕业生」,最多 200 字 diff --git a/app/service/company_service.py b/app/service/company_service.py new file mode 100644 index 0000000..8d0a948 --- /dev/null +++ b/app/service/company_service.py @@ -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 diff --git a/app/service/recruit_announcement_service.py b/app/service/recruit_announcement_service.py index 76658fa..3fe93b2 100644 --- a/app/service/recruit_announcement_service.py +++ b/app/service/recruit_announcement_service.py @@ -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) diff --git a/requirements.txt b/requirements.txt index 8ad7f54..348054b 100644 --- a/requirements.txt +++ b/requirements.txt @@ -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