本系列博客基于一个真实可运行的电商评论舆情分析项目。
上一篇:《第4篇:LangGraph 实战:把 5 个 Agent 阶段编排成一个可观测 Pipeline》
一、痛点开场
模型分析之前,数据质量决定一切。
电商评论来自不同平台,每家的导出格式都不一样:
- 列名:
评论内容、内容、comment、content……叫什么的都有 - 编码:有的 UTF-8 带 BOM,有的是 Windows 的 GBK
- 评分:数字、
5分、好评、★★★★★ - 重复:同一批数据被运营导了三次,每次导入都会重复入库
如果不处理这些问题,模型再强也是"垃圾进垃圾出"。
这篇讲数据接入层的工程细节:解析 → 归一化 → 脱敏 → 指纹去重 → 幂等落库。
二、阶段一:解析 CSV/Excel
编码自动探测
CSV 依次尝试utf-8-sig(带 BOM)、gbk(Windows 中文)、utf-8:
def_read_csv(path)->tuple[list[dict],str]:encodings=["utf-8-sig","gbk","utf-8"]last_err=Noneforencinencodings:try:withopen(path,"r",encoding=enc,newline="")asf:reader=csv.DictReader(f)ifnotreader.fieldnames:raiseValueError("CSV 无表头")returnlist(reader),encexcept(UnicodeDecodeError,ValueError)ase:last_err=eraiseValueError(f"CSV 解析失败(尝试 utf-8/gbk 均失败):{last_err}")Excel 读取
用openpyxl只读模式,首行当表头,跳过全空行:
def_read_excel(path)->list[dict]:wb=openpyxl.load_workbook(path,read_only=True,data_only=True)ws=wb.active headers=next(ws.iter_rows(values_only=True))rows=[dict(zip(headers,r))forrinws.iter_rows(values_only=True)ifrisnotNoneandany(cisnotNoneandstr(c).strip()forcinr)]wb.close()returnrows三、阶段二:字段归一化
列别名映射
_COLUMN_ALIASES把五花八门的表头映射到 canonical 字段:
_COLUMN_ALIASES={"评论内容":"content","内容":"content","评论":"content","正文":"content","评价内容":"content","comment":"content","content":"content","商品":"product_name","产品":"product_name","product_name":"product_name","平台":"platform_code","platform":"platform_code",..."用户名":"user_name","昵称":"user_name","username":"user_name",..."评分":"user_rating","星级":"user_rating","rating":"user_rating",..."评论时间":"review_time","时间":"review_time","time":"review_time",...}匹配前先归一化表头(去空白、转小写):
def_normalize_header(h)->str:return_WS_RE.sub("",str(hor"")).strip().lower()def_match_column(header)->Optional[str]:norm=_normalize_header(header)return_COLUMN_ALIASES.get(norm)or_COLUMN_ALIASES.get(norm.lower())评分容错
支持数字、5分、好评/中评/差评,最后 clamp 到 0~5:
def_parse_rating(value)->Optional[float]:# "5" → 5.0 | "5分" → 5.0 | "好评" → 5.0 | "中评" → 3.0 | "差评" → 1.0# 解析失败返回 None,行计入错误统计...returnmax(0.0,min(5.0,rating))时间容错
支持十几种格式 + 时间戳:
_TIME_FORMATS=["%Y-%m-%d %H:%M:%S","%Y-%m-%d %H:%M","%Y-%m-%d","%Y/%m/%d %H:%M:%S","%Y/%m/%d","%Y年%m月%d日","%Y.%m.%d",...]四、阶段三:用户名脱敏
隐私是硬约束,_mask_name只保留首字符:
def_mask_name(name)->str:ifnotname:return""returnstr(name)[0]+"***"生产还需要对手机号、邮箱、地址做 PII 检测与脱敏,原始文件设访问权限和保留期限。
五、阶段四:指纹去重(核心)
指纹怎么算
def_build_fingerprint(tenant_id,platform_code,product_name,content)->str:key="|".join([str(tenant_id),str(platform_codeor""),str(product_nameor""),_normalize_text(content),# 去空白,保证格式差异不影响幂等])returnhashlib.sha256(key.encode("utf-8")).hexdigest()为什么要包含这四个维度:
- 租户:不同品牌之间不冲突
- 平台 + 商品:同文案跨商品不误判
- 正文:内容本体
两层去重
# 第一层:内存 set,去当前批次内重复,减少无效写入seen=set()forrowinvalid_rows:fp=_build_fingerprint(...)iffpinseen:duplicate_count+=1continueseen.add(fp)# 第二层:数据库唯一约束,跨批次/并发/重启后的最终保障# 建表:# UNIQUE (tenant_id, fingerprint)# 插入:# INSERT ... ON CONFLICT (tenant_id, fingerprint) DO NOTHING真正的幂等不能只靠 Python set,必须依赖数据库唯一约束。因为:
- 跨批次重复:第一批导过,第二批再来
- 并发请求:两个请求同时导入
- 进程重启:内存 set 清零
六、阶段五:幂等落库
asyncdefsave_batch_node(state)->dict:# 1. 解析/创建 platform、product 维度# 2. 创建 import_batches 批次记录# 3. 批量 INSERT ... ON CONFLICT DO NOTHING# 4. 回写批次统计:inserted / duplicate / failed批次统计从落库事实汇总,而不是只信内存计数。
七、坏行处理:行级隔离
一行数据错误,不应该让整批失败:
# 缺正文 → 该行计入 failed_rows,记录原因# 评分解析失败 → 评分置空,行照常入库# 文件损坏 / 无表头 → 批次级错误,整批失败当前实现只保留前 20 条解析错误,生产应提供错误报告文件供用户下载修正。
八、踩坑与边界
- 评分/时间解析失败要保留统计:不能静默丢弃,否则运营不知道数据质量有问题
- 指纹误判:同一评论被用户修改、同文案不同人发,都可能误判重复;有外部
review_id时应优先用业务唯一键 - 大文件:当前整文件读入内存,生产要流式解析 + 分片落库
ON CONFLICT DO NOTHING的统计口径:要区分"当前批次内重复"和"历史已存在",分开统计更清晰
九、总结
- 解析兼容:编码、格式、表头全兜住
- 归一化:别名映射 + 评分/时间容错
- 脱敏:隐私前置处理
- 幂等:内存预去重 + 数据库唯一约束双保险
- 可观测:批次统计反映真实落库结果
下一篇预告:《FastAPI 长任务异步化实践:从"HTTP 阻塞"到进程内 TaskManager》