本文档拆解 dw-platform「数据提取」模块的两条核心链路:规则管理如何把圈选规则变成一次提数任务(含排队、双执行路径、回调与重试),以及人群包从创建、生成、去重到推送的完整生命周期。所有节点均标注对应的服务模块与类。执行方(datapick-service / Azkaban)的内部细节见《提数执行方详解》。
系统是 Feign 分层的 SOA 架构。Web 层(admin / h5)不连数据库,全部透传到核心服务 dw-platform-service;真正的提数执行方有两个——自建提数服务 dw-platform-datapick-service(直连 MySQL)和外部 Azkaban 大数据平台,按规则的 quotaCode 是否命中 dataPick.quotaCodes 配置分流。
flowchart LR
subgraph FE["前端 dw-platform-admin (Vue2)"]
A1["规则管理页
ruleManage.vue"]
A2["人群包管理页
packSearching.vue"]
end
subgraph WEB["Web 层(无数据库,Feign 透传)"]
B1["admin
/api/ruleManager
/api/packageManager
/api/packageSchedule"]
B2["h5
/api/package/callback"]
end
subgraph SVC["核心服务 dw-platform-service"]
C1["RuleManagerServiceImpl
规则 CRUD · generate()"]
C2["PickDataManager
排队 · 分流 · 重试"]
C3["PackageManagerServiceImpl
人群包 CRUD · 推送 · 切包"]
C4["ScheduleTaskManager
定时任务调度"]
end
subgraph EXEC["提数执行方"]
D1["datapick-service
自建提数(直连 MySQL)"]
D2["Azkaban 大数据平台
executeFlow"]
end
E[("FTP 文件服务器")]
A1 --> B1
A2 --> B1
B1 --> C1 --> C2
B1 --> C3
C4 --> C1
C2 -->|"quotaCode 命中
dataPick.quotaCodes"| D1
C2 -->|"其他数据源"| D2
D1 --> E
D2 --> E
D1 -->|"HTTP 回调"| B2
D2 -->|"HTTP 回调"| B2
B2 --> C3
| 服务模块 | 角色 | 关键入口 |
|---|---|---|
dw-platform-web/dw-platform-admin | 管理后台 Web 层 | RuleManagerWebApi、PackageManagerWebApi、PackageScheduleWebApi |
dw-platform-web/dw-platform-h5 | 回调接收端 | /api/package/callback(PackageWebApi,异步转发) |
dw-platform-soa/dw-platform-service | 核心业务:规则、提数编排、人群包、定时调度 | /service/ruleManager、/service/packageManager |
dw-platform-soa/dw-platform-datapick-service | 自建提数:MySQL 取数、去重、CSV、FTP 上传 | /service/dataPick/execute |
| Azkaban(外部) | 大数据平台提数 | POST /executor(executeFlow) |
规则存在 MySQL 表 callout_package_rule(Entity:CalloutPackageRuleDOWithBLOBs,大字段单独存放)。规则的「圈选条件」核心是 rule_condition 这个 JSON 字段,多条规则之间通过并 / 交 / 差组合,最终由提数服务翻译成一条查询。
flowchart LR R["callout_package_rule
一条提数规则"] R --> B["基础信息
rule_code · rule_name · operator
push_type · is_use"] R --> C["圈选条件 rule_condition (JSON)
按 conditionType 分四种形态"] R --> D["提数控制
quota_code 数据源
select_fields 输出字段
pre_condition · has_head"] R --> E["排除与去重
exclude 固定剔除源
exclude_condition 自定义排除
tag_exclude_condition
distinct_fields · sort_fields"] C --> C1["normal:条件数组
UNION / INTER / DIFF 组合"] C --> C2["sql:直接 SQL"] C --> C3["scala:代码片段"] C --> C4["user_domain:UserQuery
logic + 嵌套 conditions"]
rule_condition 按 condition_type 分四种形态,normal 型是最常见的一种。normal 型是一个条件数组,每个元素指明与前一规则的组合方式、数据源和一组「运算符 → 字段 → 值」的条件:
flowchart TD J["rule_condition: [ ... ]"] J --> E1["元素 1
type: UNION / INTER / DIFF
quotaCode: 数据源编码"] J --> E2["元素 2 ..."] E1 --> CM["conditionMap
{ 运算符 : { 字段 : 值 } }"] CM --> O1["Equal / NotEqual"] CM --> O2["In / NotIn"] CM --> O3["GreaterThan(LessThan)(Equal)"] CM --> O4["Like / NotLike / Between"] CM --> O5["IsNull / NotNull"] CM --> O6["fileurl
运营商匹配文件 FTP 路径"]
前端在 newRuleDetail.vue(2335–2380 行)拼装 conditionMap 后提交;用户域查询则使用 UserQuery 的嵌套结构(logic: and|or,条件算子如 geq / leq / isIn / containsArray / between)。
主入口是规则管理页的「生成人群包」按钮 → POST /api/ruleManager/generate。核心服务先做校验与排队,再按数据源分流到两条执行路径之一。下面这张是全文档最重要的一张图:
flowchart TD S["前端 ruleManage.vue
点击「生成人群包」"] -->|"POST /api/ruleManager/generate"| W["RuleManagerWebApi.generate"] W -->|"Feign /service/ruleManager/generate"| A["RuleManagerApiImpl.generate"] A --> G["RuleManagerServiceImpl.generate()"] subgraph CHECK["校验与准备"] G --> G1{"校验
packageName / operator / ruleId"} G1 --> G2{"包名唯一性
getByPackName"} G2 -->|重名| FAIL0["返回失败"] G2 -->|通过| G3["查规则 callout_package_rule
构建 PickDataConditionBO
packageCode = SqlUtil.generateCode(pack)"] end G3 --> P["PickDataManager.pickExecutor(bo, true)"] P --> Q{"acquireProject(quotaCode)
向 Redis 申请提数 project"} Q -->|"无可用资源"| WQ["加入等待队列
PICK_TASK_WAITING_QUEUE:{flow}
(Redis ZSet)
pullStatus = 1 队列中"] Q -->|"申请成功"| EXE["进入执行队列
PICK_TASK_EXECUTING_QUEUE
pullStatus = 2 提取中"] EXE --> BR{"quotaCode 是否命中
dataPick.quotaCodes 配置"} BR -->|"命中(如 user_info)"| DP["路径 1:自建提数
executeByDataPickService
Feign → datapick-service"] BR -->|"未命中"| AZ["路径 2:大数据平台
executeByBigDataService"] subgraph PATH1["datapick-service 内部(全程异步)"] DP --> D1["DataPickServiceImpl.execute
MySQLQueryBuilder 拼 SQL"] D1 --> D2["分页拉取 fetchData
敏感字段 AesUtil 解密"] D2 --> D3["dataFilter 文件匹配过滤"] D3 --> D4["deduplicate 布隆过滤器去重"] D4 --> D5["writeToCSV 生成结果文件"] D5 --> D6["FTP 上传 + HTTP 回调
callbackToThirdParty"] end subgraph PATH2["Azkaban 大数据平台"] AZ --> A1["login() 获取 session.id"] A1 --> A2["getNormalParam
conditionMap → 类 SQL 文本"] A2 --> A3["POST /executor executeFlow
flowOverride extra_param = Base64(JSON)"] A3 --> A4["返回 execid 记入 Redis
平台异步执行"] end D6 --> CB["POST /api/package/callback (h5)"] A4 -.->|"执行完成回调"| CB WQ -.->|"*/5s taskCycleExecutor
扫描队列释放资源后补位"| EXE
quotaCode:命中自建配置走 datapick-service 直连 MySQL,否则走 Azkaban。两条路径最终都通过回调汇合。PickDataManager.pickExecutor 位于 618–666 行,路径分派位于 122–132 行。除手动生成外还有两个并发入口会汇入同一条编排链路:
RuleManagerServiceImpl.add()(158–165 行)在填写了 packageName 时保存后立即调用 generate;/api/packageSchedule/generate 写入 callout_package_schedule_task,由 ScheduleTaskManager 每 3 秒扫描 Redis ZSet dw-platform-service:scheduleTask 到期触发(详见第 06 节)。提数执行完成后(无论哪条路径),文件已上传 FTP,执行方回调 /api/package/callback(h5 的 PackageWebApi,用 CompletableFuture.runAsync 异步转 Feign 到 service)。核心服务在 PackageManagerApiImpl.callback()(140–215 行)里做一串编排,人群包的 pullStatus 随之流转。
stateDiagram-v2 direction LR [*] --> Q1 : 创建 · 等待资源 Q1 : 1 队列中 QUEUEING Q2 : 2 提取中 PULLING S3 : 3 提取完成 PULL_SUCCESS F4 : 4 提取失败 PULL_FAIL D5 : 5 去重中 DISTINCTING D6 : 6 去重完成 DISTINCT_FINISH Q1 --> Q2 : acquireProject 成功 Q2 --> S3 : 回调成功 Q2 --> F4 : 回调失败 / 2h 超时 F4 --> Q2 : retryPick 重试(最多 3 次) S3 --> D5 : 配置了标签去重 D5 --> D6 : 去重完成
PackPullStatusEnum)。注意失败重试有条件限制,且超时阈值固定 2 小时。flowchart TD CB["POST /api/package/callback
→ PackageManagerApiImpl.callback()"] CB --> R1["returnApplyingProject
从 Redis 执行队列移除"] R1 --> R2{"失败?"} R2 -->|是| R3["retryPick 判断重试
(remark 含「重试」且次数 < 上限)"] R3 --> R4["AlarmSendManager 失败告警"] R2 -->|否| R5["标签去重
PackageDistinctManager.distinctByTag"] R5 --> R6["人群包打标
DataPackageManager.packageTag"] R6 --> R7["checkAndSavePackage(异步)
下载校验 FTP 文件
更新 pullStatus=3 + ftpPath + peopleNumber"] R7 --> R8["processCutPackage
切包逻辑处理"]
checkAndSavePackage 位于 PackageManagerServiceImpl 1775–1840 行。人群包是提数的产物(表 callout_package_task),通过 rule_code 与规则关联。创建入口有六个,生成之后还有去重、切包、推送、转赠、打标等一整套运营操作。
flowchart TD
subgraph CREATE["① 六个创建入口"]
E1["规则手动生成(主入口)
/api/ruleManager/generate"]
E2["新建规则时直接生成
add() 带 packageName"]
E3["定时任务生成
callout_package_schedule_task"]
E4["上传人群包
/saveByUpload · /uploadPackage"]
E5["接收外部推送
/service/packageManager/receive
大数据推送直接落库为已完成"]
E6["外部系统批量生成
outGenerate()
运营商匹配 + 临时数据替换切分"]
end
E1 --> GEN
E2 --> GEN
E3 --> GEN
E4 --> DB
E5 --> DB
E6 --> GEN
GEN["② 提数编排与执行
(见第 03 / 04 节)"]
DB[("callout_package_task
人群包记录")]
GEN --> DB
subgraph OPS["③ 人群包运营操作(packSearching.vue)"]
direction LR
O1["组合去重
distinctPackage"]
O2["人群包切分
cutPackage"]
O3["人群包合并
mergePackage"]
O4["批量打标 / 转赠 / 剔除"]
end
DB --> OPS
subgraph PUSH["④ 推送(PushSettings.vue → setCrowd)"]
direction LR
P1["按渠道推送
pushByQuotaCode"]
P2["延时推送
delayedPushTime"]
P3["推送重试
pushSchedule */1min"]
P4["结果消费
RabbitMQ push:resultQueue"]
end
OPS --> PUSH
PUSH --> DOWN["下游触达:营销 / 短链 / 评分等"]
callout_package_task 的 pullStatus / pushStatus 驱动。stateDiagram-v2 direction LR [*] --> U0 : 创建 U0 : 0 未推送 U1 : 1 推送中 U2 : 2 推送成功 U3 : 3 推送失败 U4 : 4 定时推送 U5 : 5 推送结束 U0 --> U1 : setCrowd 发起推送 U0 --> U4 : 设置延时推送 U4 --> U1 : 到达 delayedPushTime U1 --> U2 : 渠道成功 U1 --> U3 : 渠道失败 U3 --> U1 : pushSchedule 重试 U2 --> U5 : 推送结束
PackPushStatusEnum)。推送明细由 RabbitMQ 队列 dwPlatform:push:resultQueue 消费入库 callout_package_push_detail。callout_package_task.rule_code ↔ callout_package_rule.rule_code;is_use=1)时禁止删除(RuleManagerServiceImpl 254–267 行);push_type(manual / daily / week / month)决定是否支持定时生成。整条链路的可靠性依赖几个横切机制,理解它们才能解释人群包为什么会在某个状态停留:
等待队列 PICK_TASK_WAITING_QUEUE:{flow} 与执行队列 PICK_TASK_EXECUTING_QUEUE 均为 ZSet,按数据源分队列。资源有限时任务排队(pullStatus=1),释放后补位。
taskCycleExecutor 每 5 秒扫描执行中队列,超过 2 小时的任务判失败(pullStatus=4),并处理排队任务的补位调度。
retryPick 仅当回调 remark 含「重试」字样且次数未超上限才重新入队,最多 3 次,指数退避 60×(2n−1) 秒。
ScheduleTaskManager 扫描 Redis ZSet dw-platform-service:scheduleTask,到期任务调 DataPackageManager.scheduleExecutor 走同一套提数编排。
提数层:datapick-service 内布隆过滤器去重(DeduplicationManager);人群包层:PackageDistinctManager 标签去重与组合去重(*/5s 队列调度)。
推送明细走 RabbitMQ dwPlatform:push:resultQueue 异步入库;推送中的包由 pushSchedule 每分钟扫描重试。
| 层次 | 文件 | 职责 |
|---|---|---|
| 前端 | views/dataExtraction/ruleManage.vue | 规则列表、生成人群包(418 行)、生成定时任务(450 行) |
views/dataExtraction/newRuleDetail.vue | 规则编辑,conditionMap 拼装(2335–2380 行) | |
views/dataExtraction/packSearching.vue + PushSettings.vue 等 | 人群包运营操作与推送设置 | |
| 核心服务 | RuleManagerServiceImpl.java | 规则 CRUD、generate(270–332 行)、outGenerate |
PickDataManager.java | 排队(618–666 行)、分流(122–132 行)、Azkaban 对接、超时守护(672 行)、重试(1258–1306 行) | |
PackageManagerServiceImpl.java | 人群包 CRUD、callback(590–617 行)、checkAndSavePackage(1775–1840 行) | |
| 提数服务 | DataPickServiceImpl.java + MySQLQueryBuilder.java | SQL 拼装、分页取数、布隆去重、CSV 生成、FTP 上传与回调 |
| 调度 / 推送 | ScheduleTaskManager.java、PackagePushManager.java、PackageDistinctManager.java | 定时调度、推送监控重试、去重队列调度 |