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

人群包拆包详解

「拆包」在本代码库中有两种语义:业务切包——把一个整包按规则(数量 / 城市 / 机型 / 任意列 / 多列组合)切成多个子人群包,核心引擎是 PackageCutManager.isNewPackageLine,有三个触发入口(运营手动、精准系统回调联动、定时任务推送时);物理拆分——把大 CSV 按大小(90MB)或行数(500 万行)拆成多个分片上传,服务于推送渠道与撞库平台。链路总览见《规则提数与人群包链路》,提数执行细节见《提数执行方详解》

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

10两种拆包语义:业务切包 vs 物理拆分

代码里没有统一的「拆包」抽象,两个不同层面的动作都被称为拆包 / 切包,先区分清楚:

业务切包(本页主线)物理拆分
目的把整包人群拆成多个子人群包(各自有包名、数量、FTP 路径、可独立推送)把一个大文件拆成多个分片文件(合起来还是同一批人群)
产物callout_package_task 新增记录(subcontract_type=2 分包,带 parent_package_code临时文件 xxx_1.csv / xxx_2.csv…,用完即删
核心代码PackageCutManager.isNewPackageLine(行级判定)+ PackageManagerServiceImpl.asyncCutPackageCsvUtil.splitFile(按字节)+ BatchComputeHelper.splitFileByLine(按行数)
拆分依据数量上限 + 城市归属地 / 机型 / 任意列值 / 多列组合文件大小阈值(90MB)/ 行数阈值(500 万行)
触发方运营手动、精准系统回调联动、定时任务推送媒体平台推送渠道、友盟撞库数据集上传
flowchart TD
  subgraph BIZ["业务切包:整包 → 多个子人群包"]
    W[("整包 callout_package_task
subcontract_type=1")] W -->|"isNewPackageLine 判定每行归属"| N1["子包A(新记录)"] W -->|"不满足条件的行留在整包"| W2["整包(原地更新数量与ftpPath)"] end subgraph PHY["物理拆分:大文件 → 多个分片"] F["大 CSV 文件"] -->|"splitFile 按字节 / splitFileByLine 按行"| P1["分片1(含表头)"] F --> P2["分片2(含表头)"] F --> P3["分片N(含表头)"] end
图 13 · 两种拆包的产物差异。业务切包会「削」整包(留下的行写回新整包文件),物理拆分只做文件切割、每个分片都带表头。

11切包的三个入口

同一个切分引擎 PackageCutManager,被三处复用,但触发链路与持久化方式完全不同:

flowchart TD
  subgraph E1["入口① 运营手动(packSearching 页面)"]
    A1["POST /service/packageManager/cutPackage
PackageManagerApiImpl:300"] A1 --> A2["checkCutPackageParam 参数校验
(提数完成/非推送中/数量/名称唯一)"] A2 --> A3["Redis 锁 CUT_PACKAGE_KEY:{id}
TTL 10min,防并发"] A3 --> A4["异步 asyncServiceExecutor
返回「异步切包中」"] end subgraph E2["入口② 精准系统回调联动(outGenerate)"] B1["POST /service/ruleManager/outGenerate
外部精准系统按数量订购人群"] B1 --> B2["generate() 发起提数
+ 落 push_market_record(cutInfo)"] B2 --> B3["提数回调 /api/package/callback
checkAndSavePackage 完成"] B3 --> B4["processCutPackage
按实际人数重新分配切包数量"] B4 --> B5["cutPackage() → 走入口① 同一条路
成功后 pushToMarket 推回精准"] end subgraph E3["入口③ 定时任务推送时切包(TaskSettingsPage 配置)"] C1["任务设置页配置「人群包切包」
pushPackage + newPackages 存入 schedule ext"] C1 --> C2["ScheduleTaskManager 到期触发
DataPickScheduleServiceImpl.pushDataPackageSchedule"] C2 --> C3["pushToChannel → cutPackage
带 PushPackageSnapshot 快照缓存"] end A4 --> ENG["PackageCutManager.isNewPackageLine
行级判定引擎"] B5 --> ENG C3 --> ENG2["同一引擎(DataPickScheduleServiceImpl
自带的 cutPackage 重载,临时文件式)"]
图 14 · 三个入口汇入同一个行级判定引擎。入口①②落库为新人群包记录;入口③只在 FTP 上生成新文件 + 更新推送明细,不新增 callout_package_task 记录。

入口② 的数量重新分配(processCutPackage:619)

精准系统订购时填写的是「期望数量」,但实际提出来的人数以回调为准,所以回调时按 min(期望数量, 剩余总数) 顺序分配:分包名自动生成为 原包名-序号-数量。分配失败(总数为 0)会反向 HTTP 通知精准(/api/asset/pickCallback)。只有一份订购单且数量 ≥ 实际总数时跳过切包直接整包推送。

12切分引擎:isNewPackageLine 的判定顺序

引擎对每一行数据做一次判定:返回 true 写入新包,false 留在整包。判定有严格的短路顺序(PackageCutManager:16):

flowchart TD
  L["line = CSV 数据行, header = 表头
newCount = 新包已累计行数"] L --> Q1{"newPackNum ≠ null 且
newCount ≥ newPackNum?"} Q1 -->|"是"| OLD["false → 留整包
(数量上限优先于一切条件)"] Q1 -->|"否"| Q2{"cutType 为空?"} Q2 -->|"是"| NEW["true → 进新包
(纯按数量:顺序取前 N 行)"] Q2 -->|"否"| Q3{"cutType = ?"} Q3 -->|"1 城市"| C1["手机号列取前7位
查归属地表 mobile_city.csv
城市 contains 任一 typeValues"] Q3 -->|"2 机型"| C2["phone_type / 手机机型列
值 contains 任一 typeValues"] Q3 -->|"3 空号结果"| C3["改写 cutType 为校验结果列名
按列值 contains 切"] Q3 -->|"5 多列组合"| C5["multiColumnFilters 全部满足(AND)
age_range 走数值区间
其余列值 contains 任一 values"] Q3 -->|"其他(含 4 按列)"| C4["cutType 即列名
列值 contains 任一 typeValues"] C1 --> R{"匹配?"} C2 --> R C3 --> R C4 --> R C5 --> R R -->|"是"| NEW R -->|"否"| OLD
图 15 · isNewPackageLine 判定决策树。注意三点:①数量上限先于条件判定(配了 cutType 也会被数量截断);②cutType 为空等价于「顺序取前 N 条」;③所有列匹配都是 contains 模糊包含而非相等。

五种切分方式(CutType 枚举 + 前端可选)

cutType方式判定依据前端状态
(空)按数量顺序取前 newPackNum 行切分方式不选即此模式
1按手机号归属地切分手机号前 7 位 → mobile_city.csv 静态表查城市可选
2按机型切分phone_type / 手机机型 列 contains已注释隐藏
3按空号校验结果切分校验结果列值 contains(会副作用改写 cutType)已注释隐藏
4按列切分cutType 直接作为列名,列值 contains可选(支持城市/省份/自定义列下拉)
5按多列组合切分多列 AND 组合,支持 age_range 数值区间可选(多列筛选条件 UI)

列定位的隐式契约:城市切分依赖 DataSourceUtils.getMobileColumn 在表头里找到明文手机号列(phone / mobile / contact_number / handle_number / 手机号 / 联系号码… 共 11 个别名);多列切分依赖 getColumnIndex 精确匹配列名(忽略大小写)。人群包若是加密手机号列(phone_md5 等),城市切分会直接失效或取前 7 位错位(见风险清单)。

13手动切包全流程:asyncCutPackage 的「削苹果」模型

入口①②最终都走 PackageManagerServiceImpl.asyncCutPackage(:1507)。多个分包时 syncCutPackage 串行逐个执行,每切一刀整包都被「削」掉一部分,且每刀前重新查库取整包最新状态:

flowchart TD
  S["syncCutPackage
for 每个 newPackage 串行执行
每轮重新 selectByPrimaryKey 取整包"] S --> P1["1. addPackage 新增分包记录
packageCode=pack+随机码
parentPackageCode / subcontractType=2
pullStatus=2 提取中"] P1 --> P2["2. 下载整包到 ftp.cutPath
文件名 = packageCode-UUID.csv"] P2 --> P3["3. 行级分流(本地三文件并行读写)
整包读 → isNewPackageLine?
true→newWriter / false→oldWriter
表头写入两个文件,AtomicLong 计数"] P3 --> P4["4. 上传新包 → 上传削后整包
(新包失败则整体失败)"] P4 --> P5["5. 删除 FTP 原整包"] P5 --> P6["6. 更新分包:ftpPath + PULL_SUCCESS + newCount"] P6 --> P7["7. recordDO 回写 cutPackageCodes
+ pushToMarket 推送精准(MARKET 渠道)"] P7 --> P8["8. 更新整包:peopleNumber=oldCount
ftpPath 指向削后文件"] P8 --> P9["9. 两个包都上传资源服务器
(downloadUrl 供前端下载)"] P4 -.->|"任一步异常"| FAIL["失败分支:
AlarmSendManager.sendCutAlarm 告警
分包标 PULL_FAIL + remark
删除本地 old/new 文件
整包记录只写 remark 不回滚数量"]
图 16 · asyncCutPackage 主流程。整个方法在 asyncServiceExecutor 线程池里跑,无超时控制;成功后整包记录的 peopleNumber/ftpPath 原地更新(整包还是那条记录,只是「变小了」)。

并发与状态防护

14定时任务切包:快照缓存与临时文件式更新

入口③(DataPickScheduleServiceImpl:1276)与手动切包共用 isNewPackageLine,但持久化完全不同——不新增人群包记录,只生成 FTP 文件 + 更新推送明细(callout_package_schedule_task_detail):

flowchart TD
  T["pushDataPackageSchedule
解析 schedule ext 里的 pushChannelRequests"] T --> LOOP["for 每个渠道 channelRequest"] LOOP --> SNAP{"快照缓存中有
该 newPackName?"} SNAP -->|"有(本轮已切过)"| REUSE["直接复用快照
ftpPath + peopleNumber"] SNAP -->|"无"| CUT["cutPackage(newPackage, 整包ftpPath,
新UUID.csv, oldCount, newCount)
本地 .tmp 削整包 → 重命名覆盖 → 上传"] CUT --> SAVE["存入 cutPackageSnapshot
(同一配置多个渠道只切一次)"] SAVE --> UPD["推送明细记录快照的
ftpPath + peopleNumber"] REUSE --> UPD UPD --> PUSH["下载该 ftpPath 文件
走渠道推送(校验/触达)"]
图 17 · 定时任务切包(DataPickScheduleServiceImpl:1276-1400)。与手动版的差异:①整包用 .tmp 重命名式更新而非另起 UUID 文件名;②靠内存快照避免同一次调度里重复切同一配置;③无 Redis 锁、无告警、无分包落库。

配置从哪来:任务设置页(TaskSettingsPage.vue)的「人群包切包」卡片维护 cutPackageList,推送设置里每个渠道通过 pushPackage 引用某个切包配置;提交时 typeValues 从逗号字符串还原为数组,随 pushChannelRequests 存入定时任务 ext 字段(JSON)。页面会提示「先切分数据包,再校验号码」——因为切包发生在推送链路的最前面。

15物理拆分:面向渠道与撞库平台的文件切割

推送渠道拆分:CsvUtil.splitFile(按字节,90MB)

媒体对接平台渠道(PushMediumPlatform.pushToChannel,:63)在把人群包推给 CrowdPackApi.crowdPackSyn 前,先按 90MB 阈值拆分:

撞库平台拆分:BatchComputeHelper.splitFileByLine(按行数,500 万行)

友盟撞库(collision-service)上传数据集前的预处理(:189):先 Files.lines().count() 数总行数,超过 500 万 行则按 ceil(总行数/分片数) 逐行切到 .temp 工作目录(源文件名-序号.csv),逐个上传 OSS 后清空临时目录。注意此切法不写表头、按等行数切,与 CsvUtil.splitFile 语义不同。

flowchart LR
  subgraph MED["媒体平台渠道(PushMediumPlatform)"]
    M1["buildPackageAndSaveResult
校验通过的数据写成 UUID.csv"] --> M2["splitFile 90MB
每片带表头"] M2 --> M3["循环:上传资源服务器
→ crowdPackSyn 同步"] end subgraph YM["友盟撞库(BatchComputeHelper)"] Y1["本地 CSV 目录"] --> Y2["行数 > 500万?
splitFileByLine 等行数切"] Y2 --> Y3["逐片上传 OSS 数据集"] end
图 18 · 两处物理拆分。均为临时产物,不影响人群包记录与 FTP 上的整包文件。

16拆包风险清单(源码精读发现)

#位置问题影响
1PackageCutManager.cutByCity:66split[mobileColumn].substring(0, 7) 未校验长度;且加密手机号列(phone_md5 等)不在 mobileColumn 别名表内短值直接越界异常;加密包切城市时取前 7 位错乱(非异常即错切)
2cutByCity / cutByMobileType / cutByColumnType越界判断写成 split.length < index(应为 <=),matchFilter 里却是正确的 <=split.length == index 时数组越界;同库两种写法印证是笔误
3所有列匹配逻辑line.split(",", -1) 解析 CSV,未处理引号包裹的逗号(RFC 4180)含逗号字段导致列错位,切错列(上游 writeToCSV 把英文逗号替换为中文逗号,缓解但不彻底——自定义上传的包无此清洗)
4PackageManagerServiceImpl:722切包 Redis 锁 TTL 固定 10 分钟,任务无超时控制大包切分超 10 分钟后锁自动释放,可被并发切包/去重操作穿透
5asyncCutPackage 失败分支整包记录只写 remark,不回滚 peopleNumber / ftpPath;本地已上传的新 FTP 分包文件不清理失败后整包状态与 FTP 实际文件可能不一致;FTP 残留孤儿分包文件
6syncCutPackage 串行循环多个分包串行执行,每刀都重新下载/上传整包N 个分包 = N 次整包下载上传,大包场景 IO 放大 N 倍
7BLACK_NUM_RESULT 分支newPackage.setCutType(...) 修改入参对象(副作用)同一 NewPackage 对象被入口③的快照复用时,第二次切包走错分支
8processCutPackage:642分包名自动生成 包名-序号-数量,若与既有包重名将直接抛异常回调链路切包失败,精准侧只能靠告警感知
9checkCutPackageParam(批量版:804)newPackNum 允许为 null(按 0 累计),但引擎会把 null 视为「无数量上限」全部切走校验与引擎对 null 语义不一致:null 数量的分包理论上可吃掉整个整包
10DataPickScheduleServiceImpl.cutPackage定时切包无 Redis 锁、无告警;快照是 JVM 内存 Map多实例部署时同一调度任务并发触发会重复切包;失败静默,仅日志
11CsvUtil.splitFile:563异常时 catch 后返回空 List(调用方当「拆分失败」处理),掩盖真实异常渠道推送失败原因不可观测,需翻服务端日志