DW-PLATFORM · 数据提取模块

规则提数与人群包链路

本文档拆解 dw-platform「数据提取」模块的两条核心链路:规则管理如何把圈选规则变成一次提数任务(含排队、双执行路径、回调与重试),以及人群包从创建、生成、去重到推送的完整生命周期。所有节点均标注对应的服务模块与类。执行方(datapick-service / Azkaban)的内部细节见《提数执行方详解》

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

01总体链路:一次提数经过的所有服务

系统是 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
图 1 · 提数总体链路。前端 → Web 层 → Feign → 核心服务编排 → 按数据源分流到自建提数服务或 Azkaban,产物上传 FTP 后回调 h5 入口更新人群包状态。
服务模块角色关键入口
dw-platform-web/dw-platform-admin管理后台 Web 层RuleManagerWebApiPackageManagerWebApiPackageScheduleWebApi
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)

02规则数据模型:一条规则长什么样

规则存在 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"]
图 2 · 规则的数据结构。圈选条件 rule_conditioncondition_type 分四种形态,normal 型是最常见的一种。

normal 型条件的 JSON 结构

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 路径"]
图 3 · normal 型条件。多个元素按 type 依次并 / 交 / 差,最终合并成一次提数查询。

前端在 newRuleDetail.vue(2335–2380 行)拼装 conditionMap 后提交;用户域查询则使用 UserQuery 的嵌套结构(logic: and|or,条件算子如 geq / leq / isIn / containsArray / between)。

03提数主流程:从点击「生成人群包」到任务执行

主入口是规则管理页的「生成人群包」按钮 → 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
图 4 · 提数主流程。核心分流点在 quotaCode:命中自建配置走 datapick-service 直连 MySQL,否则走 Azkaban。两条路径最终都通过回调汇合。PickDataManager.pickExecutor 位于 618–666 行,路径分派位于 122–132 行。

除手动生成外还有两个并发入口会汇入同一条编排链路:

04回调处理与状态流转

提数执行完成后(无论哪条路径),文件已上传 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 : 去重完成
图 5 · pullStatus 状态机(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
切包逻辑处理"]
图 6 · 回调编排顺序。checkAndSavePackage 位于 PackageManagerServiceImpl 1775–1840 行。

05人群包管理:从创建到推送的完整生命周期

人群包是提数的产物(表 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["下游触达:营销 / 短链 / 评分等"]
图 7 · 人群包生命周期四阶段。①②③④ 之间通过 callout_package_task 的 pullStatus / pushStatus 驱动。

推送状态机(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 : 推送结束
图 8 · pushStatus 状态机(PackPushStatusEnum)。推送明细由 RabbitMQ 队列 dwPlatform:push:resultQueue 消费入库 callout_package_push_detail

规则与人群包的关联约束

06关键机制:队列、定时、重试与去重

整条链路的可靠性依赖几个横切机制,理解它们才能解释人群包为什么会在某个状态停留:

Redis ZSet 提数排队

等待队列 PICK_TASK_WAITING_QUEUE:{flow} 与执行队列 PICK_TASK_EXECUTING_QUEUE 均为 ZSet,按数据源分队列。资源有限时任务排队(pullStatus=1),释放后补位。

*/5s 超时守护

taskCycleExecutor 每 5 秒扫描执行中队列,超过 2 小时的任务判失败(pullStatus=4),并处理排队任务的补位调度。

失败重试(有条件)

retryPick 仅当回调 remark 含「重试」字样且次数未超上限才重新入队,最多 3 次,指数退避 60×(2n−1) 秒。

*/3s 定时任务调度

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.javaSQL 拼装、分页取数、布隆去重、CSV 生成、FTP 上传与回调
调度 / 推送ScheduleTaskManager.javaPackagePushManager.javaPackageDistinctManager.java定时调度、推送监控重试、去重队列调度