--- title: "03-spider_data.py" created: 2026-04-03 tags: - 项目 aliases: - spider_data.py --- # spider_data.py ### `spider_data.py` — 数据处理与上传 核心类 `SpiderData`,消费爬虫产出的 JSON+HTML 文件,驱动后续处理流水线。 `process_announcement_data()` **是核心方法**,根据 `proc_type` 参数执行不同操作: | proc\_type | 操作 | | --- | --- | | `ann` | HTML → Markdown → 大模型解析公告结构化数据 | | `wx` | 微信文章特殊处理(含图片复杂度检测)后同上 | | `cjob` | 解析公司职位数据(`parse_cjob`) | | `up_api` | 将模型结果上传到云端(学校公告) | | `up_api_cjob` | 将模型结果上传到云端(公司职位) | **其他功能:** - 黑名单系统:基于 URL、图片 MD5、文本内容 MD5 过滤垃圾数据 - 微信文章图片复杂度检测:GIF、小尺寸无文字图片过多时跳过,避免大模型处理失败 - 手工微信文章支持:`pre_wx_article()` 将本地 HTML 文件整理成统一数据格式 - 用 `.ok` / `.err` / `.expired` 后缀文件做状态标记,避免重复处理 ## 代码 ```python # -*- coding: utf-8 -*- import os import glob import json import shutil import time from utils import ner_logger,is_wechat_url,getMD5Str from utils_html import get_weixin_info,get_weixin_hand_url,clean_weixin_html from utils_date import check_file_modification_time,check_file_modification_time_old from utils_playwright import get_wx_url_content from parsegpt.ann_md import html2md_with_fix,md_to_html from parsegpt.ann_model import parse_announcement from parsegpt.cjob_model import parse_cjob from api.quanzhi_api import upload_cloud,upload_cloud_job from spider_sch import DEFAULT_COMMON # 默认手工维护微信文章的特殊学校代码(固定ID) DEFAULT_WX_SCHOOL = "sch_88888" # 默认一次处理的数据条数(防止一次性处理过多) DEFAULT_PCOUNT = 3 # 数据处理核心类:负责爬虫结果清洗、MD转HTML、大模型解析、上传云端 class SpiderData(): """数据处理中心:清洗爬虫数据 → 转换格式 → 大模型抽取 → 上传API""" def __init__(self,_spider_sch): """ 初始化数据处理器 :param _spider_sch: 爬虫主类实例(传递配置、工具方法) """ self.spider_sch = _spider_sch # ===================== 加载图片/文本黑名单(过滤垃圾内容) ===================== # 图片URL黑名单 with open("data/black_img_urls.txt",encoding="utf-8") as f: self.black_img_urls_list = f.read().splitlines() # 图片MD5黑名单 with open("data/black_img_md5.txt",encoding="utf-8") as f: self.black_img_md5_list = f.read().splitlines() # 图片内容MD5黑名单 with open("data/black_img_md5_content.txt",encoding="utf-8") as f: self.black_img_md5_content_list = f.read().splitlines() # 文本黑名单(暂时空) self.black_text_list = [] # 加载完整URL黑名单并计算MD5加入黑名单 with open("data/black_img_urls.full.txt",encoding="utf-8") as f: for line in f.read().splitlines(): _md5 = getMD5Str(line.strip()) self.black_img_urls_list.append(_md5) # 去重 self.black_img_urls_list = list(set(self.black_img_urls_list)) # 打印黑名单加载数量 ner_logger.info(f"读取黑名单数量,url:{len(self.black_img_urls_list)} ,md5: {len(self.black_img_md5_list)},text : {len(self.black_text_list)}") # ===================== 黑名单检查工具 ===================== def check_url_in_blacklist(self,_md5): """检查图片MD5是否在黑名单中""" if _md5 in self.black_img_urls_list: return True if _md5 in self.black_img_md5_list: return True return False def check_md5_txt_in_blacklist(self,_txt_md5): """检查文本内容MD5是否在黑名单""" if _txt_md5 in self.black_img_md5_content_list: return True return False # ===================== 路径/配置获取工具 ===================== def get_wx_path(self): """获取跨平台微信文章本地存放路径""" if os.name == "nt": return self.spider_sch.config.get(DEFAULT_COMMON,"wxpath_win") if os.name == "posix": return self.spider_sch.config.get(DEFAULT_COMMON,"wxpath_mac") return "" def get_wx_url(self,_data): """ 从数据中智能提取微信文章URL 优先:type_url=wxwz → 其次last_url """ if 'type_url' in _data and _data['type_url'] == "wxwz" and is_wechat_url(_data['full_url']): return _data['full_url'] elif "last_url" in _data and is_wechat_url(_data['last_url']): return _data['last_url'] else: return "" # ===================== 外部命令调用工具 ===================== def proc_wechaturl_md(self,_url,_dir,_outfile): """调用外部exe将微信文章转成MD文件""" _exe = self.spider_sch.get_md_exe() _cmd = f"{_exe} -image=url --dir {_dir} --output {_outfile} \"{_url}\" " ner_logger.info(f"开始执行微信文章转换md {_cmd}") os.system(_cmd) def proc_html_md(self,_inputfile,_outfile): """调用外部exe将HTML转成MD文件""" _exe = self.spider_sch.get_html_md_exe() _cmd = f"{_exe} --output-overwrite --plugin-table --exclude-selector=\".ad\" --input \"{_inputfile}\" --output \"{_outfile}\"" ner_logger.info(f"开始执行html转换md {_cmd}") os.system(_cmd) # ===================== 核心:处理公告/招聘数据 ===================== def process_announcement_data(self,_key,sch_info,_stat,proc_type = "ann"): """ 统一数据处理入口 :param _key: 机构ID(sch_xx / com_xx) :param sch_info: 机构配置信息 :param _stat: 全局状态(计数、去重、配置) :param proc_type: 处理类型 wx=微信 / ann=公告 / cjob=企业职位 / up_api=上传 """ _keydir = self.spider_sch.get_key_dir(_key) _ar_dir = self.spider_sch.get_savepath(f"/data/ardata/{_key}") _cache_dir = self.spider_sch.get_savepath(f"/data/cache") # 创建输出目录 if not os.path.exists(_ar_dir): os.makedirs(_ar_dir) if not os.path.exists(_cache_dir): os.makedirs(_cache_dir) ner_logger.info(f"_keydir:{_keydir},{proc_type}") # 遍历该机构下所有JSON数据文件 for _file in glob.glob(f"{_keydir}/*.json"): _filename = os.path.basename(_file) # 跳过列表页index文件,只处理详情detail_文件 if _filename.startswith("index"): continue if not (_filename.startswith("detail_") and _filename.endswith(".json")): continue # ===================== 打开JSON & 基础过滤 ===================== with open(_file,"r",encoding="utf-8") as f: _hfile = _file.replace(".json",".html") # 公告/职位必须有HTML文件 if proc_type in ["ann","cjob"] and not os.path.exists(_hfile) : continue # 标记:本次处理过的文件,避免重复 _tag_info = f"{_filename}_{proc_type}" # 10分钟内新文件 → 跳过(防止爬虫还在写入) if check_file_modification_time(_file): ner_logger.info(f"文件10s内生成,跳过 {_file}") continue # 10天前旧文件 → 跳过 if check_file_modification_time_old(_file): if not _tag_info in _stat['all_proc_list']: _stat['all_proc_list'].append(_tag_info) continue # ===================== 定义各类输出文件路径 ===================== _ar_file = os.path.join(_ar_dir,_filename) # 归档JSON _md_file = _ar_file.replace(".json",".md") # MD文件 _fix_file = _ar_file.replace(".json",".html") # 清洗后HTML _model_file = _ar_file.replace(".json",".model.json") # 大模型结果 _up_api_file = _model_file.replace(".json",f".json.{_stat['dist']}.ok") # 上传成功标记 _up_api_err_file = _model_file.replace(".json",f".json.{_stat['dist']}.err") # 上传失败 _expired_file = _ar_file.replace(".json",".json.expired") # 过期不处理 _job_ann_file = _ar_file.replace(".json",".json.job") # 职位标记 # 已过期 → 跳过 if os.path.exists(_expired_file): continue # 上传模式必须有模型文件 if proc_type in ["up_api","up_api_cjob"] and not os.path.exists(_model_file): continue # 加载原始数据 _data = json.load(f) _title = _data['announcement_name'] # 标题不符合关键词 → 跳过 if proc_type in ["wx",'ann'] and not self.spider_sch.is_title_include(_title): ner_logger.info(f"标题在关键词排除之列,{_file} / {_title}不需要爬取,跳过") continue # 获取微信URL wx_url = self.get_wx_url(_data) # 调试:只处理指定前缀文件 if 'pfile' in _stat and not _filename.startswith(_stat['pfile']): continue # 调试:强制重新处理 if 'pfile' in _stat and _filename.startswith(_stat['pfile']): if os.path.exists(_model_file) and proc_type in ["wx",'ann',"cjob"]: os.remove(_model_file) if os.path.exists(_up_api_file) and proc_type in ["up_api","up_api_cjob"]: os.remove(_up_api_file) # 已生成模型 → 跳过 if os.path.exists(_model_file) and proc_type in ["wx","ann",'cjob']: continue # 已上传 → 跳过 if os.path.exists(_up_api_file) and proc_type in ["up_api","up_api_cjob"]: continue # 本次已处理 → 跳过 if _tag_info in _stat['all_proc_list'] and not 'pfile' in _stat: _stat['all_proc_list'].append(_tag_info) continue time.sleep(0.5) # ===================== 分类型处理 ===================== # 1. 处理微信文章 if (proc_type == "wx" and wx_url) or (proc_type == "ann" and wx_url): _ok,_wx_file = self.pre_wx_article_html(_data,_cache_dir,wx_url) if _ok: _ok,_hfile,_mdfile = self.process_wechat_data(sch_info,_file,_ar_dir,_filename,_data,_cache_dir,wx_url,_wx_file) if _ok: _ok = self.process_gonggao_data(sch_info,_filename,_data,_cache_dir,_md_file,_fix_file,_model_file,_hfile,proc_type,_expired_file,_stat) if _ok == "ok": _stat['total'] = _stat.get('total',0) +1 _stat['p_count'] +=1 if _stat['total'] >= DEFAULT_PCOUNT: ner_logger.info(f"共处理{_stat['p_count']}条wx数据完成,退出") return # 2. 处理普通公告 elif proc_type == "ann" and not wx_url: _ok = self.process_gonggao_data(sch_info,_filename,_data,_cache_dir,_md_file,_fix_file,_model_file,_hfile,proc_type,_expired_file,_stat) if _ok == "ok": _stat['total'] = _stat.get('total',0)+1 _stat['p_count'] +=1 if _stat['total'] >= DEFAULT_PCOUNT: ner_logger.info(f"共处理{_stat['p_count']}条ann数据完成,退出") return # 3. 处理企业职位 elif proc_type == "cjob": ner_logger.info(f"开始处理到公司职位文件:{_data['announcement_name']} - {_hfile}") _ok,_msg = parse_cjob(self,_model_file,_data,sch_info,_expired_file,_hfile,_stat) if _ok == "ok": _stat['total'] = _stat.get('total',0)+1 _stat['p_count'] +=1 if _stat['total'] >= DEFAULT_PCOUNT: ner_logger.info(f"共处理{_stat['p_count']}条cjob数据完成,退出") return # 4. 上传公告到API elif proc_type == "up_api" and os.path.exists(_model_file) and not os.path.exists(_job_ann_file): if os.path.exists(_up_api_file): continue ner_logger.info(f"开始处理数据上传云端 - {_hfile} \n {_model_file}") _ok,_msg,_code= upload_cloud(_model_file,_stat['dist']) if _ok and _code == "200": with open(_up_api_file,"w",encoding="utf-8") as f: f.write(_msg) time.sleep(2) elif _ok and _code != "200": with open(_expired_file,"w",encoding="utf-8") as f: f.write(_msg) else: with open(_up_api_err_file,"w",encoding="utf-8") as f: f.write(_msg) # 5. 上传职位到API elif (proc_type == "up_api_cjob" or (proc_type == "up_api" and os.path.exists(_job_ann_file))) and os.path.exists(_model_file): if os.path.exists(_up_api_file): continue ner_logger.info(f"开始处理职位数据上传云端 - {_hfile} \n {_model_file}") _ok,_msg,_code= upload_cloud_job(_model_file,_stat['dist'],_stat['retry']) if _ok: with open(_up_api_file,"w",encoding="utf-8") as f: f.write(_msg) time.sleep(2) else: with open(_up_api_err_file,"w",encoding="utf-8") as f: f.write(_msg) # ===================== 公告处理:HTML → MD → 大模型抽取 ===================== def process_gonggao_data(self,sch_info,_filename,_data,_cache_dir,_md_file,_fix_file,_model_file,_hfile,proc_type,_expired_file,_stat): """处理学校/企业公告:清洗HTML → 转MD → 大模型抽取结构化数据""" ner_logger.info(f"分析文件{_data['announcement_name']} - {_md_file}") with open(_hfile,"r",encoding="utf-8") as f: _html = f.read() # HTML 转 MD + 清洗 _ok,_info,_full_text = html2md_with_fix(self,_data,_fix_file,_md_file,_html,sch_info,_cache_dir,_hfile) if _ok and len(_info) > 2 and 'props' in _info: # 微信文章额外检查图片复杂度 if proc_type == "wx": _ok = self.check_wx_file(_data,_hfile,_model_file) if not _ok: with open(_expired_file,'a',encoding='utf-8') as fw: fw.write(f"不能处理这个微信公告(283)\n") return "" # 大模型抽取结构化字段 (_ok,msg) = parse_announcement(_md_file,_fix_file,_model_file,_info,_full_text,proc_type,sch_info,_expired_file,_stat) if _ok: ner_logger.info(f"大模型处理文件成功:{_data['announcement_name']} - {_md_file}") return "ok" else: with open(_expired_file,'a',encoding='utf-8') as fw: fw.write(f"{msg}\n") else: with open(_expired_file,'a',encoding='utf-8') as fw: fw.write(f"不能处理这个微信公告(299) html2md_with_fix \n") return "" # ===================== 微信文章专用处理 ===================== def process_wechat_data(self,sch_info,_file,_ar_dir,_filename,_data,_cache_dir,wx_url,_wx_file): """微信文章专用流程:转MD → 转回HTML → 供后续清洗""" ner_logger.info(f"处理微信数据{_data['announcement_name']} - {wx_url} \n {_file}") _md_filename = _filename.replace(".json",".md") _md_cache_dir = f'{_cache_dir}_md' if not os.path.exists(_md_cache_dir): os.makedirs(_md_cache_dir) _md_file = os.path.join(_md_cache_dir,_md_filename) # 重新生成 if os.path.exists(_md_file): os.remove(_md_file) # 调用工具转MD self.proc_wechaturl_md(_wx_file,_md_cache_dir,_md_filename) if not os.path.exists(_md_file): return "","","" # MD内容过短 → 无效 with open(_md_file,"r",encoding="utf-8") as f: _md_text = f.read() if len(_md_text) < 50: return "","","" # MD 转 HTML _html_file = _md_file.replace(".md",".html") _ok = md_to_html(_md_file,_html_file) if _ok : return "ok",_html_file,_md_file return "","","" def pre_wx_article_html(self,_data,_cache_dir,wx_url): """预先下载微信文章到本地,清洗HTML,提取信息""" _wx_cache_dir = f'{_cache_dir}_wx' if not os.path.exists(_wx_cache_dir): os.makedirs(_wx_cache_dir) _md5_url = getMD5Str(wx_url) _wx_file = os.path.join(_wx_cache_dir,f"{_md5_url}.html") _wx_file_0 = os.path.join(_wx_cache_dir,f"{_md5_url}.html.0") _wx_file_config = os.path.join(_wx_cache_dir,f"{_md5_url}.html.config") # 未缓存 → 下载 if not os.path.exists(_wx_file) or not os.path.exists(_wx_file_config): executable_path = self.spider_sch.get_browser_path() _ok,_content,image_lists = get_wx_url_content(executable_path,wx_url) if _ok: with open(_wx_file,"w",encoding="utf-8") as f: f.write(_content) with open(_wx_file_config,"w",encoding="utf-8") as f: json.dump(image_lists,f,ensure_ascii=False,indent=4) else: return False,"" # 备份原始文件 if not os.path.exists(_wx_file_0): shutil.copy(_wx_file,_wx_file_0) # 提取微信文章信息(作者、时间等) get_weixin_info(_wx_file,_data) # 清洗微信HTML垃圾 clean_weixin_html(_wx_file,_wx_file_0) return True,_wx_file # ===================== 手工微信文章导入 ===================== def pre_wx_article(self): """导入本地手工整理好的微信HTML文章,生成标准JSON结构""" _wx_path = self.get_wx_path() ner_logger.info(f"wx文章路径{_wx_path}") # 固定导入到默认微信机构ID _keydir = self.spider_sch.get_key_dir(DEFAULT_WX_SCHOOL) if not os.path.exists(_keydir): os.makedirs(_keydir) # 遍历所有HTML for _file in glob.glob(f"{_wx_path}/**/*.html",recursive=True): _file_err_file = _file.replace(".html",".html.err") _filename = os.path.basename(_file) # 提取标题、链接、日期 title,link, wxdate = get_weixin_hand_url(_file) _file_json = f'detail_{getMD5Str(link)}.json' _file_file = f'{_keydir}/{_file_json}' # 标记不可处理 → 删除并跳过 if os.path.exists(_file_err_file): if os.path.exists(_file_file): os.remove(_file_file) continue # 已处理 → 跳过 if os.path.exists(_file_file): continue # 构造标准数据结构 _data = { "announcement_name":title, "publish_time":wxdate, "link":link, "full_url":link, "last_url":"", "parent_url":"http://mp.weixin.qq.com", "upload": "", "contact": "无", "type_url":"wxwz", "channel": DEFAULT_WX_SCHOOL, "wx_name":"", "wx_title":title, "wx_public_time":wxdate, "wx_source_file":_file, } # 写入标准JSON with open(_file_file,"w",encoding="utf-8") as f: f.write(json.dumps(_data,ensure_ascii=False,indent=4)) # ===================== 微信文章质量检查(过滤复杂图片) ===================== def check_wx_file(self,_data,_hfile,_model_file): """ 检查微信文章是否适合自动处理 规则:小图/动态图过多 → 判定复杂,不处理 """ count = 0 _props = _data['props'] img_urls = _props.get('img_urls',{}) wx_source_file = _hfile if 'wx_source_file' in _data: wx_source_file = _data['wx_source_file'] _ok_file = wx_source_file.replace(".html",".html.ok") _err_file = wx_source_file.replace(".html",".html.err") # 图片复杂度评分 for _url,_dict in img_urls.items(): img_ocr = _dict['img_ocr'] img_width = _dict['img_width'] img_height = _dict['img_height'] if 'full_qr' in _dict and _dict['full_qr'] == 'Y': continue # 极小图 if img_width * img_height < 100 * 100 and len(img_ocr) < 5: count += 2 # 中小图 elif img_width * img_height < 500 * 500 and len(img_ocr) < 5: count += 1 # GIF动图 elif _url.endswith(".gif"): count += 2 # 分数过高 → 不处理 if count > 10: with open(_err_file,"w",encoding="utf-8") as f: f.write("不好处理") return False # 标记可处理 with open(_ok_file,"w",encoding="utf-8") as f: f.write(f"可以处理{_model_file}") return True ``` --- **项目分区导航**:[[02-spider_com.py|spider_com.py]] ⬅️ | 03-spider_data.py | ➡️ [[04-spider_sch.py|spider_sch.py]]