从一套乱七八糟的多源用户数据开始说起吧。上个月刚接手一个活儿,运营那边丢了三份名单过来,客服一份、订单一份、埋点一份,每份名单里都有“用户ID”,但三份名单的ID体系完全对不上。同一人在我面前顶着三个ID、两个手机号、三种昵称,再怎么join都拼不出一张干净的大宽表。做数据集成的人应该都经历过这种时刻——主键靠不住,只能靠算法把“看起来是同一个人”的记录找出来。这篇是数据集成系列的第二篇,专门讲算法层的东西:实体识别、相似度计算、阻塞策略、匹配判定、流式增量整合,以及我实际跑完三版之后踩出来的坑。适合正在做数仓、数据中台、实时数仓,天天被多源数据折磨的工程师和数据科学方向的朋友参考。
1. 数据集成算法到底在解决什么问题
1.1 从“同一个人,三个ID”说起:先把问题边界摸清楚
数据集成不是简单的join。传统join的前提是两条记录共享同一个可靠的键,哪怕是代理键也行。但真实的多源数据里,很多场景下不存在这样一个全局键:A系统的用户ID是自增整数,B系统用的是手机号掩码,C系统干脆只有设备ID和昵称。你没法用任何一个字段硬join,只能综合多个字段判断“这大概率是同一个人”。
这个判断过程,学术界叫记录链接(Record Linkage),工业界叫实体解析(Entity Resolution),本质上是一个决策问题:给定来自不同系统的两条记录,基于它们各属性的相似程度,决定是否指向同一真实世界实体。
举个当时项目里的例子。客服名单里的记录是:
user_id=1001, 手机=138****1234, 姓名=张伟订单名单里的记录却长这样:
customer_id=A56B, 手机=13812345678, 收货人=张伟峰, 地址=北京市海淀区中关村大街1号埋点名单里更疯:
device_id=MX-2024-3345, 昵称=weizhang, 设备=小米14, 地址=中关村大街1号院这三条记录没有共同主键,但手机号一致、姓名高度相似、地址基本同一。数据集成算法要做的,就是把这些特征综合起来,输出“同一个人”的判断,并给一个置信度。所以你看,它解决的问题不是“怎么把数据搬过来”,而是“搬过来之后怎么判断它们原本是同一个东西”。
1.2 分清一层关系:数据集成算法不是ETL,也不是join
很多同学一上来就把数据集成和ETL混在一起说,这里我尽量给你划条线。
ETL关注的是数据的同步、清洗、格式转换,主力工具是DataX、SeaTunnel、DolphinScheduler这类调度同步平台。它们解决的是“数据从A系统到B系统,保证量、保证速、基本保证质”。数据集成算法关注的是数据的内容语义:两条记录描述的是不是同一个对象、同一笔事件、同一个产品。
join则是SQL层面的概念,它的前提是键已经明确。数据集成算法处理的是“键不明确”的情况。你可以这么理解:ETL把原材料搬进仓库,join是把仓库里有明确料号的东西合并出库,数据集成算法则是在料号贴错、缺失、重复的情况下,靠规格和经验把东西认出来。
这个边界不搞清楚,做技术方案时经常会出现一种尴尬:你跟业务说“我们做用户主数据融合”,业务以为你要把所有表join一遍,结果交付物却是一套“相似度打分+人工审核”的机制。所以动手前先把定位对齐,后面才少吵架。
2. 实体识别:整个匹配算法的主战场
2.1 相似度计算:从编辑距离到向量距离
实体识别的第一步,是计算字段间的相似度。这一层选错,后面全盘皆输。我见过不少刚入门的朋友一提到相似度就想到Levenshtein编辑距离,然后拿它算地址、算商品名、算一切,结果效果很差还不知道为什么。
不同字段类型要选不同的相似度度量,这是基础中的基础。
常用相似度算法大致可以分几个梯队:
| 算法 | 适合场景 | 特点 |
|---|---|---|
| Levenshtein编辑距离 | 短字符串、拼写错误 | 直观,但对换序和长文本不友好 |
| Jaro-Winkler | 人名、较短名称 | 前缀权重高,比较吃字段序 |
| Jaccard + n-gram | 地址、短语、有缩写场景 | 对语序变化有一定容忍度 |
| SimHash / MinHash | 长文本、大规模近似去重 | 能压缩、能分桶,适合海量数据 |
| 词向量 / 句子向量 | 跨语言、同义替换 | 需要训练或预训练模型,成本较高 |
拿地址来说,“北京市海淀区中关村大街1号”和“中关村大街1号,北京市海淀区”语义完全一致,但逐字符编辑距离远得离谱。我的做法是先做地址归一化:省市区分开、去除行政后缀、统一路名简称,然后再用Jaccard或者字符n-gram相似度,效果比直接编辑距离好一个量级。
若要用Python快速验证相似度,常见组合是:
import numpy as np from scipy.spatial.distance import pdist, squareform from itertools import combinations def char_ngrams(s, n=2): return set(s[i:i+n] for i in range(len(s)-n+1)) def jaccard(a, b): if not a or not b: return 0.0 s1, s2 = char_ngrams(a), char_ngrams(b) return len(s1 & s2) / len(s1 | s2) records = ["北京市海淀区中关村大街1号", "中关村大街1号,北京市海淀区", "中关村南大街1号"] mat = squareform(pdist(np.array(records).reshape(-1, 1).astype(object), lambda x, y: 1 - jaccard(x[0], y[0]))) print(mat)但说句实话,字段相似度再精密,也只是整个判定链路里的一个输入。更关键的一步,是要先做预处理:大小写、全半角、简繁体、电话号码区号补齐、地址的省市区拆解。没有这一步,再强的算法都是拿脏数据做精细运算,属于事倍功半。
2.2 阻塞策略:把全两两比较变成分桶比较
实体识别理论上要对所有记录做两两比较。但一百万条记录做两两比较,就是5千亿个记录对,单机跑一天都出不来,更别提生产环境要天天跑。所以工程上必须做阻塞(Blocking)——先用一个廉价但有效的条件,把所有记录分到若干桶里,只在桶内做候选对比较。
阻塞策略选得好不好,直接决定算法的召回率。如果阻塞字段本身有错,千辛万苦调的相似度模型根本连候选对都没见过,那就是典型的“在源头漏掉”。
几种常用的阻塞思路:
- 规则分桶:比如按“省份+姓氏拼音首字母+性别”分桶。代价低,但桶可能严重倾斜——北京、王姓这种组合可能几百万条挤一个桶。
- 排序键分桶:先把字段转成soundex、拼音首字母、数字组合等排序键,再按排序键区间分段。对姓名的容错能力不错。
- LSH分桶:用MinHash + LSH对长文本做近似匹配,把相似文档映射到同桶。适合地址、长文本、商品描述这类数据。
- 枚举阻塞:如果同一实体可能在数据里以多种方式出现,比如手机号有可能带“+86”也有可能不带,可以把多种规范化形式都建索引,条条大路通罗马。
我自己的经验是,现实项目里必须组合多重阻塞,别指望一个钥匙开所有锁。宁可让候选对多一点,靠后面的精细相似度过滤,也不要阻塞太紧导致漏匹配。因为这个取舍,如果你把精确率做上去了但召回率很低,业务方查不到人的时候,回来骂的还是你。
2.3 匹配判定:规则阈值组合怎么定
有了候选对和字段相似度,下一步是判定“到底像到什么程度才算同一人”。最常见的做法是加权得分。我习惯给每个证据维度打分:
- 手机号完全相同:+0.8
- 姓名相似度大于0.85:+0.15
- 地址相似度大于0.7:+0.05
- 总分 ≥ 0.9:判定为同一实体
权重怎么调?不是拍脑袋。我的做法是抽500条候选对,人工标注“是同人/不是同人”,然后根据这500条的得分分布来定阈值。画一张查准率-查全率曲线,看业务对误配和漏配的容忍度:如果有客服人工介入环节,可以放宽阈值追求高召回;如果不需要人工,系统自动做批量合并,那就得紧一点,误配代价太高。
得分结果最好分三档:高置信自动合并、中置信人工确认、低置信直接丢弃。这个分档机制比一个单一阈值实用得多,因为真实数据的分布不是你阈值一刀切就能切干净的。中置信这批人,往往才是你下一轮优化算法最有价值的样本。
3. 把匹配问题交给机器学习
3.1 特征工程与样本标注:先有评判标准,再谈模型
规则调阈值也有天花板。字段一多、脏数据模式一多,写规则的人会疯掉。这时候就该上机器学习方案了。
前期工作其实不复杂:把候选对转成特征向量。特征包括:各属性的相似度得分、字段缺失情况、来源系统、时间差、渠道组合交叉特征。这里要特别提醒一句,负样本别随手从全量数据里随机抽。我从实战总结出来的经验是,负样本分布要和真实阻塞后候选对的分布一致,否则模型在训练时学到的概率分布,和上线时面对的分布完全对不上,效果必然崩塌。
正样本从哪里来?两种主流方案:
- 用现有规则模型的高置信结果作为正样本,再人工抽检清洗。
- 从低置信、中置信区域抽样让人工标注,这样可以教模型学到规则之外的模式。
第二种就是我们常说的主动学习。主动学习的核心不是让模型替你干活,而是让模型告诉你它“最不确定”的那些样本在哪,让有限的标注人力花在刀刃上。
3.2 分类模型与主动学习的实操套路
把匹配判定当成二分类问题来做,常用模型是逻辑回归和随机森林。随机森林在这种场景下表现相当稳,对离散特征和缺失值容忍度高,也不容易被个别极端特征带偏。热搜词里有一个“随机森林回归算法”,那是回归问题用的,匹配判定是分类问题,别混用模型,原理可以共用但任务目标要清晰。
我一般会跑一个这样的流程:
- 先跑一个简单规则模型,输出高置信度匹配对,纳入训练集。
- 对规则模型拿不准的区域做不确定性采样,人工标注500到1000对。
- 把两类样本合并,训练随机森林分类器。
- 在保留测试集上评估查准率、查全率、F1值。
- 持续迭代:把每次人工审核的结果回灌到训练集。
这里有个常见的坑,就是类别不平衡。匹配问题里负样本数量远大于正样本。一个候选对里真正“同人”的比例可能不到1%。解决方式不是简单过采样或下采样,而是要看模型输出的概率分数是否可信。如果需要精确概率值,记得对预测概率做校准,可以用Platt缩放,否则你画出来的阈值曲线会失真。
3.3 实体聚类:把两两配对变成实体簇
匹配判定的输出是“记录对是否同实体”,但真正要落到数仓里的,是一个个实体簇:用什么作为统一的实体ID,簇里包含哪些系统里的哪些记录,置信度是多少。这个过程中需要聚类算法来把两两关系聚成更大的实体簇。
常见的做法有两种。一是连通分量聚类,简单粗暴;二是在连通分量的基础上加置信度限制,比如两条记录置信度高于0.95才合并。我强烈建议不要用纯连通分量。因为它在传递关系下会产生雪崩效应:A像B,B像C,结果A和C长得完全不像也被合并成一个实体。这在用户画像里是要命的问题。
一个更稳妥的办法是贪心聚类加剪枝:先把候补对按相似度从高到低排序,一条条决定是否并入已有簇,同时设置每簇的记录数上限和最大传递步数。听着朴素,但工程稳定性比一个花哨的谱聚类靠谱得多。
有少数媒体稿已经把“暴力枚举算法”吹成万金油,但在实体聚类这种量级下,暴力枚举的空间你根本吃不下,要做的是在候选空间里做剪枝。剪枝是一种策略,不是偷懒。
4. 增量与流式集成:数据不会等你全量重算
4.1 为什么全量重算跟不上节奏
如果你的用户主数据一周更新一次,离线全量重算完全没问题。但实时营销、实时风控这类对时效要求高的场景,每天一个新注册用户进来,你不可能每天晚上做一个全量实体解析,再等到第二天给他打上标签。这时候就需要增量匹配和流式集成的思路。
流式数据环境里有三个绕不开的问题:乱序、迟到、重复。你收到一条新注册事件,它对应的老用户可能在三分钟前已经更新过手机号,事件顺序颠倒了。你直接用最新事件去匹配老实体,很可能因为信息不完整而错过。这项工程比离线匹配复杂得多。
4.2 增量匹配的四种常用手段
增量匹配不是另起炉灶,而是在离线实体解析的基础上,多加一层在线索引和状态管理。这里分享四种我常用且验证过的手段。
第一种,索引复用。每个实体簇维护一个“特征摘要”,比如簇内所有手机号的布隆过滤器、所有姓名拼音的集合、地址中的高频词集合。新记录进来时,先查布隆过滤器,快速判断这个手机号是不是已经存在于某个簇里。如果存在,直接跟该簇做精细匹配;如果不存在,再去建立新候选。
第二种,LSH动态分桶。MinHash + LSH的好处在于,相似记录大概率会落在同一个哈希桶里。新记录到达时,只要查对应桶即可拿到候选,复杂度从O(n)降到常数级。代价是LSH只是近似匹配,可能漏掉边缘情况,所以要用多个LSH分桶方案融合,避免单一哈希族太偏。
第三种,时间窗口衰减。实体特征也会变老。一个人五年不活跃,他的旧地址跟新记录做相似度匹配的意义就很小。可以只对最近N天活跃的实体簇做强匹配,更久远的簇放到每周/每月的离线低频重扫。这样状态不会无限膨胀,匹配质量也更聚焦。
第四种,定期重算与增量融合。当前比较稳的架构是离线每天做一次全量实体解析,生成全局实体ID;实时链路负责当日新增数据的增量匹配,并将结果映射到离线生成的ID上。两个链路用同一套特征和规则版本,ID映射表放Redis或HBase,保证两边对得上。增量链路和离线链路之间通常会有一个对账环节,把差异拉到人工审核表里。
这里顺便说一句,DolphinScheduler这类调度平台适合编排离线任务和定时同步,它解决不了实时流的匹配问题。很多人把调度平台和流计算混为一谈,实际上实时增量匹配得靠Flink或Spark Streaming这类流处理框架自己维护状态。调度平台管的是“几点跑、依赖谁、失败重跑”,流处理管的是“每条消息来了怎么算”。
5. 实操避坑:我跑了三版才有稳定结果
5.1 数据量级决定算法选型
数据集成算法没有银弹,量级不一样,做法天差地别。我一般用一个经验表来选型。
| 数据量级 | 建议方案 | 备注 |
|---|---|---|
| 10万级 | 全两两比较 + 字符串相似度 | 单机即可,Python完全能扛 |
| 10万~100万 | 阻塞 + 规则加权 | 重点关注阻塞字段质量 |
| 100万~1000万 | 分布式阻塞 + 分类模型 | Spark或Flink,注意key倾斜 |
| 1000万以上 | LSH + 增量索引 | 在线索引和离线批处理结合 |
| 在线毫秒级 | 布隆过滤器 + 缓存 | 只做快速预判,精细匹配异步完成 |
选型时有一条原则:先看字段缺失率。缺失率超过50%的字段,不管理论多漂亮,都不要放进特征集合。它只会带来两个问题:一是候选对被无效分割,二是模型学到的是“缺失模式”而不是“真实相似度”。
5.2 高频问题速查表:我踩过的坑都在这里
跑多了以后,你会发现大部分问题其实都集中在那几个点上。整理成表格给你,直接对着现象查就行。
| 现象 | 可能原因 | 对策 |
|---|---|---|
| 匹配结果忽高忽低 | 源系统字段口径变化 | 建立固定评测集,每次发布前回放评测 |
| 误合并越来越多 | 阈值太低,或特征里放了不该有的宽字段 | 提高阈值,增加高危否定规则,比如身份证号不同直接排除 |
| 漏匹配越来越多 | 阻塞字段错误或缺失 | 换多重阻塞策略,把单字段阻塞改成组合阻塞 |
| 内存一直爆 | 阻塞桶倾斜,个别桶记录太多 | 对桶再拆分,或对桶内数据做采样近似 |
| 同一记录被并入多个簇 | 连通分量传递闭包导致 | 限制传递步数,或按置信度做贪心聚类 |
| 模型解释不了 | 随机森林特征权重漂移 | 定期输出特征重要性,和业务核对维度是否合理 |
另外还有几个非算法层面的心得,比调算法重要得多。
第一,数据使用合规。你在做实体识别的时候,会把手机、姓名、地址等字段拿来做匹配计算。这中间涉及到的数据安全和个人信息保护,一定要提前和你所在公司的安全合规团队对齐,明确哪些字段能用于匹配、哪些只能脱敏后统计、匹配结果谁有权查看。这不是技术问题,却是能让你整套方案停摆的问题,提前沟通比事后补救便宜太多。
第二,一定要做人审闭环。再先进的模型也有拿不准的时候。中置信区域拉一个审核队列,让业务方每天花15分钟过一眼,把结果回灌到训练集。这个操作带来的长期收益,比我调任何算法的效果都明显。
第三,版本管理。匹配规则或者模型版本一定要记录,每次升级前先做回归。我以前吃过一次亏:模型新版本上线,某类商品实体识别率涨了一个点,但另一类商品召回率直接掉了六个点,因为新旧数据分布微妙变化,没有评测集兜底,差点上线翻车。现在我把评测集当成核心资产保存,所有改动都必须跑回归。
6. 工具生态:算法引擎选型与造轮子边界
6.1 调度和传输层:算法往往不在这一层
很多人一搜“数据集成”,看到DolphinScheduler、DataX、SeaTunnel就以为这就是数据集成算法,这是一个很大的误区。DolphinScheduler确实可以编排数据集成任务,比如每天凌晨2点同步订单数据、3点跑全量实体解析、4点对账,但它本身不提供实体匹配的算法。DataX和SeaTunnel同样如此,它们是传输管道,把数据从A搬到B,内容语义的判定还得在拿到数据之后自己实现。
我刚工作那阵也犯过这个错,项目一开始就搭了一堆调度任务,以为把管道打通就完事。后来业务问“一个用户在三套系统里怎么合并”的时候,才发现真正需要投入精力的不是管道,而是下游的匹配逻辑。所以进项目之后,建议先想清楚哪些任务是同步,哪些任务是集成,哪些任务是调度,提前分配资源。
6.2 真正做实体解析的开源库
如果你的场景不是非得造轮子,有几套开源方案可以直接拿来用。
Python Record Linkage Toolkit,适合中中小规模数据,单机运行,内置了大量相似度算法,配合pandas用起来很顺手。Dedupe更偏主动学习,你提供字段,它引导你标注样本,然后训练出匹配模型,适合字段不固定、没有强主键的场景,我个人在有些用户主数据项目里用过,效果不错,但需要接受它的交互式标注流程,有时会有点费人力。Spark Entity Resolution这类适合大规模数据,做分布式阻塞和匹配,但配置成本较高,得懂Spark调优。
如果你用的是ClickHouse、StarRocks这类OLAP引擎,里面也有近似查询函数,如SimHash、Jaccard距离,适合在SQL层面做初步去重,但别指望它们能完成复杂的实体解析。SQL层面的近似查询更适合做快速过滤,精细判定还是交给专门流程处理。
6.3 自研与现成的边界
我的判断标准很简单两条:第一,你的业务实体是不是有强领域规则;第二,字段和实体形态是不是高度不稳定。
如果实体判定极度依赖领域字典,比如医疗里的药品名称归一化、法律里的判决文书案由匹配,自研规则和特征更有优势,通用的开源库很难覆盖这类语义。如果只是普通的用户实体识别,字段无非姓名、手机号、邮箱、地址,Dedupe或者Record Linkage Toolkit足够,没有必要自己从零开始写相似度和分类器。
但如果你决定自研,有一件事必须做:建立一套稳定的评测集。没有评测集的调参,本质上就是自嗨。每次改特征、改阻塞策略、换模型,都要拿到评测集上跑一把,看查准率、查全率、F1怎么变。数据集成算法的维护是长期的,评测集是这个长期过程中最可靠的锚点。
最后分享一个小技巧。做匹配的时候,高置信度的匹配结果不要只停留在“打标签”,尝试把结果回填到源系统的ID映射表。比如找到“客服系统的1001=订单系统的A56B”之后,把这条映射写回去。这样下次增量匹配时,系统可以优先信任已知映射关系,积少成多,整套集成算法会越用越轻,也越用越准。我个人在实际操作中的体会是,数据集成算法真正难的从来不是某一篇论文里的花哨模型,而是怎么把归一化、特征、阻塞、阈值、人工审核这些脏活累活串成一条稳定流水线。把这些基本功打磨好,你的方案怎么都不会跑偏。