DW-PLATFORM · DOCS · 数据提取模块

人群包去重详解

「去重」在本代码库中是三条独立链路组合去重(两包求差、A 包原地削、内存二分查找)、按类型剔除去重(规则 / 标签 / 内部包,产生新包,本地走 ExternalSortUtil 外部排序,或下发 Azkaban file[ftp] flow)、标签去重(生成时回调自动触发,借道 datapick-service 与 people_tag 比对)。链路总览见《规则提数与人群包链路》

导读 规则与人群包总览 提数执行方详解 人群包拆包详解 人群包去重详解

17三条去重链路总览

三个入口互相独立,算法、产物、状态机都不同,不要混为一谈:

flowchart TD
  subgraph A["链路 A · 组合去重(packSearching 勾两包)"]
    A1["POST /distinctPackage(List 版)
PackageManagerServiceImpl:829"] A1 --> A2["A 包标「去重中」→ 异步"] A2 --> A3["下载 A/B 两包全量进内存
B 包手机号排 long[] → 二分查找
A = A − B"] A3 --> A4["结果覆盖上传回 A 包原 FTP 路径
A 包原地削,不产生新记录"] end subgraph B["链路 B · 按类型剔除(弹窗选择)"] B1["POST /distinctPackage(BO 版)
PackageManagerServiceImpl:891"] B1 --> B2{"distinctType"} B2 -->|"RULE 规则"| B3["按规则条件重新走提数
pickExecutor → Azkaban"] B2 -->|"TAG 标签"| B4["查 package_tag 拿 FTP 列表
拼 file[ftp] 条件"] B2 -->|"PACKAGE 内部包"| B5["被选包 ftpPath 作副包
拼 file[ftp] 条件"] B4 --> B6{"useLocalDistinct?
(prod: true)"} B5 --> B6 B6 -->|"规则去重 或 false"| B7["pickExecutor
Azkaban/数综 file-ftp flow"] B6 -->|"标签/内部包 且 true"| B8["PackageDistinctManager
本地去重(外部排序)"] B7 --> B9["新包记录 createPackage
(受理后插入)"] B8 --> B9 end subgraph C["链路 C · 标签去重(生成时自动)"] C1["提数回调 /api/package/callback
成功 + 规则配了 tagExcludeCondition"] C1 --> C2["PackageDistinctManager.distinctByTag:654"] C2 --> C3["datapick-service /distinctByTag
与 people_tag 比对剔除已触达"] C3 --> C4["带新 ftpPath 再次回调
按新文件落库"] end
图 19 · 三条去重链路。A 改原包;B 产生新包(原包不动);C 在生成链路上自动剔除非用户操作。
链路 A 组合去重链路 B 按类型剔除链路 C 标签去重
触发运营勾选两个包「组合去重」运营弹窗选类型(规则/标签/内部包)生成回调时自动(规则配置驱动)
算法内存:B 排序 long[] + 二分查找本地:ExternalSortUtil 外部排序;或 Azkaban file flowdatapick-service 分批查 people_tag
产物A 包被原地削(无新记录)新人群包记录(createPackage)新 FTP 文件,覆盖回调结果
原包状态去重中(5) → 去重完成(6)不动不动
并发控制无锁(校验推送状态兜底)Redis 双 ZSet 队列 + */5s 守护(本地路径)无(回调链路串行)

18链路 A:组合去重(A − B,二分查找版)

flowchart TD
  S["distinctPackage(List<Integer>)
PackageManagerServiceImpl:829"] S --> V["校验:恰好 2 包 / 人数 > 0
/ 非推送中 / 提取完成或去重完成"] V --> M["A 包 pullStatus=5 去重中
同步 updateByPrimaryKey"] M --> ASYNC["asyncServiceExecutor 异步
asynDistinctPackage:1384"] ASYNC --> D1["下载 A、B 两包
downloadAndAnalysis 全量读入 List"] D1 --> ALG["getDistinctDataBin:1466
① B 包每行取第一列,合法手机号→long,否则→0
② parallelStream 排序成 long[]
③ 遍历 A 包:第一列非手机号→直接保留
④ 是手机号→binarySearch B 数组
找不到才保留"] ALG --> CNT["peopleNumber = 结果行数
(有表头再 -1)"] CNT --> UP["生成 CSV 覆盖上传回
A 包原 FTP 路径"] UP --> F["finally:A 包 pullStatus=6 去重完成
失败时 remark 写原因"] F --> RES["成功:重新上传资源服务器
(供前端下载)"]
图 20 · 组合去重流程。B 包完全不动,A 包被原地削掉交集;算法 O(n log m),靠 long 数组而非 HashSet 省内存。

算法细节(getDistinctDataBin:1466):排序阶段对 B 包用 parallelStream().mapToLong().sorted();过滤阶段对 A 包逐行 Arrays.binarySearch。两处都只看第一列,且用 DataSourceUtils.isPhoneNumber 判定是否明文手机号——非手机号行(包括表头)直接保留。

19链路 B:按类型剔除(产生新包)

三种 distinctType 的条件构建

类型去重依据构建条件(PackageManagerServiceImpl)
RULE 规则去重用户选一条规则file[ftp](main_file_url isIn A包路径) + exclude 规则排除项(:1193)
TAG 标签去重选标签 + 时间范围package_tag(tag + operator + 时间窗)拿 FTP 列表 → file[ftp](main_file_url…;ext_file_url…;select_fields…)(:1117)
PACKAGE 内部包去重选另一个包同上,ext_file_url = 被选包 ftpPath;去重列先固定 phone,异步下载后按表头提取(:1059)

执行路径分叉(:935)——规则去重永远走提数编排层;标签/内部包去重看开关:

新包记录的插入时机:与规则生成一致——先发去重请求,拿到 OK / queued / failed 后才 createPackage 插入(:952),初始 pullStatus 分别为「提取中(2) / 队列中(1) / 失败(4)」,此时 ftpPath、peopleNumber 均为空,等回调回填。

20本地去重引擎:队列编排 + ExternalSortUtil

编排层:PickDataManager 的「复刻」

PackageDistinctManager(:88)与提数编排层结构几乎完全对称,仅语义有两处不同:

组件PickDataManager(提数)PackageDistinctManager(去重)
等待/执行队列Redis 双 ZSet + flow 索引 Set同构,key 前缀 PackageDistinctManager:,member 为 BO JSON / packageCode
*/5s 守护循环2h 超时标失败 + FIFO 补位完全相同(:156)
acquireProjectAzkaban project 池轮转(全局 4 个)执行队列 size < maxQueueSize(生产 6)即放行(:520)
受理后动作真实 execid 入执行队列packageCode 入执行队列,回调凭它出队

执行层:DistinctPackageServiceImpl → ExternalSortUtil

flowchart TD
  E["distinctPackage(bo)
DistinctPackageServiceImpl:43
CompletableFuture.runAsync 触发"] E --> M["1. 下载主包(重试版)
去重列:配置值 或 按表头自动提取
(优先明文手机号列,其次 uu_id 列)"] M --> X["2. 逐个下载副包
单个失败仅 warn 被忽略
全部失败才整体失败"] X --> H["3. checkHeader 表头校验
主包 + 每个副包都必须包含全部去重列
任一缺失即整体失败"] H --> S1["4a. 分块排序:每 10 万行在内存排序
写成临时文件 sortedBlock_N.txt
排序 key = 去重列拼接值
(主包行携带完整内容,副包只存 key)"] S1 --> S2["4b. 多路归并去重(小顶堆)
相邻相等 key 只输出一次
→ 主包、副包各自一份有序无重文件"] S2 --> S3["4c. 双指针过滤 filterMainFile
主包行 key 在副包出现→丢弃
否则写出原始完整行 + 写表头 + 计数"] S3 --> U["5. 结果上传 FTP dataExtractPath"] U --> CB["6. 回调自己 packageManagerApi.callback
sourceService=dw-platform-service
(回调侧凭此跳过 FTP 校验)"] CB --> ST["新包状态 → 提取完成
ftpPath / peopleNumber 回填"]
图 21 · 本地去重主流程。ExternalSortUtil(:63)三步:分块排序 → 归并去重 → 双指针过滤,全程磁盘临时文件接力,内存占用被 maxLineInMemory(10 万行)封顶,专为超大包设计。

为什么链路 A 不用外部排序:组合去重(:1466)仍是全内存实现(两个 List + long 数组),而外部排序只服务链路 B 的本地路径——两条链路的算法代际不同,是历史演进的结果。

21链路 C:标签去重(生成时自动剔除已触达)

挂在提数回调链上(PackageManagerApiImpl:171),完全由规则配置驱动,用户无感知:

sequenceDiagram
  participant EXE as 提数执行方
  participant CB as /api/package/callback
  participant PDM as PackageDistinctManager
  participant DPS as datapick-service

  EXE->>CB: 提数成功回调(ftpPath=人群包)
  CB->>CB: distinctByTag 前置判断
  Note over CB: ① success=true ② 回调非 DISTINCT/PICK 类型
(防递归)③ 规则存在且配了 tagExcludeCondition CB->>PDM: distinctByTag(callback) PDM->>PDM: 标签条件(标签+时间窗)→ package_tag id 列表
buildPickRequest 注入 packageTag PDM->>DPS: POST /distinctByTag(distinctFtpPath=人群包) DPS-->>PDM: 受理返回 code=0 DPS->>DPS: 异步:下载人群包 CSV
每 1000 行一批,去重列 AES 加密后
查 people_tag(打过的=已触达) DPS->>DPS: 剔除后写新文件 → 上传 FTP DPS->>CB: 带新 ftpPath 再次回调(pickType=DISTINCT) CB->>CB: 跳过标签去重(防递归)→ 按新文件落库
图 22 · 标签去重时序。本质是「生成即自动剔除历史上打过的标签人群」,用两次回调完成,第二次回调凭 pickType=DISTINCT 跳过再去重。

22回调收尾:一条回调链服务提数与去重

/api/package/callback(PackageManagerApiImpl:140)同时承接提数回调与本地去重回调,靠两个信号分流:

pullStatus 相关状态

状态含义谁写入
5 去重中(DISTINCTING)组合去重进行中链路 A 入口同步设置(:871)
6 去重完成(DISTINCT_FINISH)组合去重结束(无论成败)链路 A finally 必写(:1449)
2/1/4链路 B 新包的受理态去重请求返回后 createPackage 前(:946)
3 提取完成链路 B 回调后终态checkAndSavePackage 回填

切包的前置校验允许「提取完成 或 去重完成」,即链路 A 完成后的包仍可继续被切包;去重与切包之间靠 DISTINCT_PACKAGE_KEY Redis 锁互斥。

23去重风险清单(源码精读发现)

#位置问题影响
1asynDistinctPackage:1404链路 A 全内存:两包整个读入 List + long[]千万级大包 OOM 风险;外部排序(链路 B)未反向覆盖链路 A
2getDistinctDataBin:1481只看第一列,且非手机号行直接保留第一列不是明文手机号的包等于没去重(表头也因此被保留,靠 isDataSourceHasHead 修正计数)
3DistinctPackageServiceImpl:99副包下载失败仅 warn 即忽略去重结果不完整且无告警,用户无感知
4asynDistinctPackage:1443上传失败分支 return false 后,异常分支缺 return(方法声明 boolean)编译器静默返回 false,异常时返回值语义不可靠
5PackageDistinctManager 超时守护2h 超时直接标新包 PULL_FAIL,但不中断异步线程本地文件照常写、回调照常来,状态先失败后又被回调改回,存在状态抖动
6链路 A 并发控制无 Redis 锁,仅校验「非推送中」两人同时对同一 A 包发起组合去重,两份结果互相覆盖(切包路径有锁,去重路径没有)
7buildPickDataConditionForDistinctPackage:1080去重列硬编码 "phone",靠 needExtractColumns 异步修正注释自述:若改回中台去重需先按表头取列,当前是权宜写法
8checkHeader:197表头包含判断用 header.contains(column) 字符串包含列名互为子串时误判通过(如 phone 与 phone_md5)
9PackageDistinctManager:51Logger 取的是 PickDataManager.class去重日志全部打在提数 logger 名下,排障易误导
10链路 B 的新包插入时机与提数共用「受理后插入」,但去重失败重试逻辑(retryPick)依赖 remark 含「重试」本地去重失败 remark 为异常消息(不含「重试」),不会自动重试
11链路 A 状态写入「去重中」用 updateByPrimaryKey 全字段更新(:873),并发读改写与其他操作(打标/切包)并发时可能互相覆盖彼此字段