「去重」在本代码库中是三条独立链路:组合去重(两包求差、A 包原地削、内存二分查找)、按类型剔除去重(规则 / 标签 / 内部包,产生新包,本地走 ExternalSortUtil 外部排序,或下发 Azkaban file[ftp] flow)、标签去重(生成时回调自动触发,借道 datapick-service 与 people_tag 比对)。链路总览见《规则提数与人群包链路》。
三个入口互相独立,算法、产物、状态机都不同,不要混为一谈:
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
| 链路 A 组合去重 | 链路 B 按类型剔除 | 链路 C 标签去重 | |
|---|---|---|---|
| 触发 | 运营勾选两个包「组合去重」 | 运营弹窗选类型(规则/标签/内部包) | 生成回调时自动(规则配置驱动) |
| 算法 | 内存:B 排序 long[] + 二分查找 | 本地:ExternalSortUtil 外部排序;或 Azkaban file flow | datapick-service 分批查 people_tag |
| 产物 | A 包被原地削(无新记录) | 新人群包记录(createPackage) | 新 FTP 文件,覆盖回调结果 |
| 原包状态 | 去重中(5) → 去重完成(6) | 不动 | 不动 |
| 并发控制 | 无锁(校验推送状态兜底) | Redis 双 ZSet 队列 + */5s 守护(本地路径) | 无(回调链路串行) |
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["成功:重新上传资源服务器
(供前端下载)"]
算法细节(getDistinctDataBin:1466):排序阶段对 B 包用 parallelStream().mapToLong().sorted();过滤阶段对 A 包逐行 Arrays.binarySearch。两处都只看第一列,且用 DataSourceUtils.isPhoneNumber 判定是否明文手机号——非手机号行(包括表头)直接保留。
| 类型 | 去重依据 | 构建条件(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)——规则去重永远走提数编排层;标签/内部包去重看开关:
packageDistinct.useLocal=false 或 RULE 类型 → pickDataManager.pickExecutor(bo, true):走 Azkaban file-ftp flow,由大数据平台读 FTP 文件完成去重;useLocal=true(生产配置)且 TAG / PACKAGE → packageDistinctManager.distinctPackage(bo, true):本服务本地去重(见下节)。新包记录的插入时机:与规则生成一致——先发去重请求,拿到 OK / queued / failed 后才 createPackage 插入(:952),初始 pullStatus 分别为「提取中(2) / 队列中(1) / 失败(4)」,此时 ftpPath、peopleNumber 均为空,等回调回填。
PackageDistinctManager(:88)与提数编排层结构几乎完全对称,仅语义有两处不同:
| 组件 | PickDataManager(提数) | PackageDistinctManager(去重) |
|---|---|---|
| 等待/执行队列 | Redis 双 ZSet + flow 索引 Set | 同构,key 前缀 PackageDistinctManager:,member 为 BO JSON / packageCode |
| */5s 守护循环 | 2h 超时标失败 + FIFO 补位 | 完全相同(:156) |
| acquireProject | Azkaban project 池轮转(全局 4 个) | 执行队列 size < maxQueueSize(生产 6)即放行(:520) |
| 受理后动作 | 真实 execid 入执行队列 | packageCode 入执行队列,回调凭它出队 |
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 回填"]
为什么链路 A 不用外部排序:组合去重(:1466)仍是全内存实现(两个 List + long 数组),而外部排序只服务链路 B 的本地路径——两条链路的算法代际不同,是历史演进的结果。
挂在提数回调链上(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: 跳过标签去重(防递归)→ 按新文件落库
/api/package/callback(PackageManagerApiImpl:140)同时承接提数回调与本地去重回调,靠两个信号分流:
pickDataManager.returnApplyingProject 出执行队列 + retryPick 判断重试;packageDistinctManager.returnApplyingProject(:588,按 packageCode 从去重执行队列移除);| 状态 | 含义 | 谁写入 |
|---|---|---|
| 5 去重中(DISTINCTING) | 组合去重进行中 | 链路 A 入口同步设置(:871) |
| 6 去重完成(DISTINCT_FINISH) | 组合去重结束(无论成败) | 链路 A finally 必写(:1449) |
| 2/1/4 | 链路 B 新包的受理态 | 去重请求返回后 createPackage 前(:946) |
| 3 提取完成 | 链路 B 回调后终态 | checkAndSavePackage 回填 |
切包的前置校验允许「提取完成 或 去重完成」,即链路 A 完成后的包仍可继续被切包;去重与切包之间靠 DISTINCT_PACKAGE_KEY Redis 锁互斥。
| # | 位置 | 问题 | 影响 |
|---|---|---|---|
| 1 | asynDistinctPackage:1404 | 链路 A 全内存:两包整个读入 List + long[] | 千万级大包 OOM 风险;外部排序(链路 B)未反向覆盖链路 A |
| 2 | getDistinctDataBin:1481 | 只看第一列,且非手机号行直接保留 | 第一列不是明文手机号的包等于没去重(表头也因此被保留,靠 isDataSourceHasHead 修正计数) |
| 3 | DistinctPackageServiceImpl:99 | 副包下载失败仅 warn 即忽略 | 去重结果不完整且无告警,用户无感知 |
| 4 | asynDistinctPackage:1443 | 上传失败分支 return false 后,异常分支缺 return(方法声明 boolean) | 编译器静默返回 false,异常时返回值语义不可靠 |
| 5 | PackageDistinctManager 超时守护 | 2h 超时直接标新包 PULL_FAIL,但不中断异步线程 | 本地文件照常写、回调照常来,状态先失败后又被回调改回,存在状态抖动 |
| 6 | 链路 A 并发控制 | 无 Redis 锁,仅校验「非推送中」 | 两人同时对同一 A 包发起组合去重,两份结果互相覆盖(切包路径有锁,去重路径没有) |
| 7 | buildPickDataConditionForDistinctPackage:1080 | 去重列硬编码 "phone",靠 needExtractColumns 异步修正 | 注释自述:若改回中台去重需先按表头取列,当前是权宜写法 |
| 8 | checkHeader:197 | 表头包含判断用 header.contains(column) 字符串包含 | 列名互为子串时误判通过(如 phone 与 phone_md5) |
| 9 | PackageDistinctManager:51 | Logger 取的是 PickDataManager.class | 去重日志全部打在提数 logger 名下,排障易误导 |
| 10 | 链路 B 的新包插入时机 | 与提数共用「受理后插入」,但去重失败重试逻辑(retryPick)依赖 remark 含「重试」 | 本地去重失败 remark 为异常消息(不含「重试」),不会自动重试 |
| 11 | 链路 A 状态写入 | 「去重中」用 updateByPrimaryKey 全字段更新(:873),并发读改写 | 与其他操作(打标/切包)并发时可能互相覆盖彼此字段 |