提数的「怎么执行」:编排层 PickDataManager 的 project 资源池与 Redis 队列模型、自建 datapick-service 的主提数循环(SQL 拼装 → 游标分页 → 解密 → 布隆去重 → CSV → FTP → 回调)、Azkaban 大数据平台的对接五步,以及源码精读发现的 11 项实现风险。链路总览见《规则提数与人群包链路》。
编排层是提数的「调度大脑」,核心概念是 project 资源池:project 是 Azkaban 上的执行项目(提数作业容器),同一 project 同时只能跑一个任务。prod 环境池子为 data_extract, data_extract_1~3(共 4 个,配置 spring.bigdata.extProject)。Redis 按 flow(= quotaCode,每个数据源一套队列)维度隔离,project 池全局共享。
编排层对上游只暴露一个入口 pickExecutor(bo)(:122),内部按 quotaCode 是否命中 dataPick.quotaCodes 配置路由到两个执行引擎之一。prod 配置 quotaCodes: user_info——即目前只有 user_info 走自建服务,其余指标(各类 TM_/TD_ 表)全走 Azkaban,主力流量仍在 Azkaban 侧。两条通道完成后都通过统一回调出口 /api/package/callback(PackageManagerApiImpl:140)更新人群包状态。
flowchart TD IN["pickExecutor(bo) 入口
PickDataManager:122
设置 callbackUrl = dw.site/api/package/callback"] IN --> RT{"quotaCode ∈
dataPick.quotaCodes?"} RT -->|"是(prod: user_info)"| CHA["路径 A:自建 datapick-service
executeByDataPickService:135
Feign DataPickApi → /execute
同步返回 code=0 即受理
execid = 本地随机 UUID
不进执行队列、不占 project"] RT -->|"否(其余数据源)"| CHB["路径 B:Azkaban 大数据平台
executeByBigDataService:237
login → buildParam → /executor
返回真实 execid
进入执行队列、占用 project"] CHA -.->|"异步跑完
FTP 上传结果"| CB["统一回调出口
/api/package/callback
PackageManagerApiImpl:140"] CHB -.->|"平台异步跑完
FTP 上传结果"| CB
编排逻辑的归属关系:编排层的排队 / project 池 / 生命周期管理,本质都是为 Azkaban 通道的并发约束服务的;datapick-service 通道自身分页拉 MySQL,不受该约束,因此绕开了整套体系:
| 编排层逻辑 | 作用 | 归属 |
|---|---|---|
acquireProject() | 从 project 池挑一个未被占用的 project(随机) | 仅 Azkaban:同一 project 的 flow 只能并发跑一个,需池化轮转 |
| 等待队列 / 执行队列(Redis ZSet) | 同 flow 占满即排队,pickExecutor(bo, queue) | 仅 Azkaban:排队针对 Azkaban 并发约束 |
executeTaskTimeOutHandle() | 2 小时超时标失败、清执行队列 | 仅 Azkaban:处理 project:packageCode:execid 格式任务 |
回调后 returnApplyingProject() | 提数完成归还 project | 仅 Azkaban |
buildPickRequest() / 标签条件注入 | 构建统一请求 BO,tagExcludeCondition 转子查询条件 | 两通道共用(请求组装层) |
统一回调出口 + retryPick() | 人群包状态回写、失败重试入等待队列 | 两通道共用(结果处理层) |
flowchart TD PE["pickExecutor(bo, needQueue=true)
PickDataManager:618"] PE -->|"BO 已带 project"| EXE["直接执行"] PE -->|"project 为空"| AP["acquireProject(flow):1043
读执行 ZSet 取已占用 project
候选 = 池 − 已占用 → 随机取一个"] AP -->|"有空闲"| EXE AP -->|"全占用"| WQ2["addPickTaskWaitingQueue:965
member=BO的JSON, score=入队时间戳
返回 code=queued"] EXE --> R{"提数请求结果"} R -->|"code=OK"| AA["addApplyingProject:1141
member=project:packageCode:execid
score=受理时间戳(超时基准)"] R -->|"msg 含 running
(flow 已在跑)"| WQ2 R -->|"failed"| F["返回 failed"]
flowchart TD C["taskCycleExecutor
*/5s(:672)"] C --> T["executeTaskTimeOutHandle
扫描执行中 ZSet 全部 flow"] C --> W["waitingTaskHandle
扫描等待队列全部 flow"] T --> T1{"now − score > 2小时?"} T1 -->|"是"| T2["查 callout_package_task
非「提取完成」→ pullStatus=4
remark=提数任务执行超时失败
从 ZSet 删除"] T1 -->|"score=null"| T3["脏数据直接删除"] W --> W1["Redisson 锁
waitingPickTaskExecuteLock:{flow}
tryLock 不等待/租约10s"] W1 -->|"拿到锁"| W2["acquireProject(flow)
无空闲则本轮放弃"] W2 --> W3["按 score 升序取最早任务
(严格 FIFO)"] W3 --> W4["从等待队列删除
递归调 pickExecutor(bo, true)"] W4 --> W5{"结果回写"} W5 -->|"OK"| W6["pullStatus=2 提取中"] W5 -->|"queued"| W7["留队列等下轮"] W5 -->|"failed"| W8["pullStatus=4 + 删除"]
| Key(dw-platform-service 域) | 结构 | member / score | 用途 |
|---|---|---|---|
pickTaskWaitingQueue:{quotaCode} | ZSet | member=BO 的 JSON / score=入队时间戳 | 等待队列,TTL 2 天 |
pickTaskExecutingQueue:{flow} | ZSet | member=project:packageCode:execid / score=受理时间戳 | 执行中任务(占用 project),超时判定基准 |
dwPlatform:retryPickTask:packageCode:{code} | String | 重试次数 | 重试计数,TTL 2 天 |
buildPickDataCondition(RuleManagerServiceImpl:341)用 copyProperties 把规则 DO 的同名字段注入 BO(sortFields / distinctFields / tagExcludeCondition / conditionType / exclude / quotaCode / selectFields / preCondition 等),rule_condition 显式映射到 condition。三类 ExecuteType 决定回调去向:人群包(package)/ 临时数据(tempPackage,文件搬入 FTP /upload/temp/,7 天生命周期)/ 定时任务数据(schedulePackage)。
dataPick.quotaCodes,prod 即 user_info)「受理即返回 + 异步执行」模型:execute 同步阶段只做校验和 SQL 构建(CompletableFuture.runAsync,默认 ForkJoinPool),立即返回 code=0;真正提数在异步线程里跑完整个主循环,结束后 HTTP 回调。编排侧 executeByDataPickService(PickDataManager:135)通过 Feign 接口 DataPickApi(dw-platform-datapick-api 模块,指向 dw-platform-datapick-service)调用,同步返回 code=0 即视为受理成功,execid 为本地随机 UUID(非真实执行 id),且不进执行队列、不占 project——即完全绕开编排层排队体系,仅复用统一回调出口。
flowchart TD E["execute(req)
DataPickServiceImpl:121"] --> QB["getQueryBuilder:359
quotaCode 直接作表名,单表无 JOIN"] QB --> BR{"quotaCode ==
dataPick.JD.table?"} BR -->|"是"| JD["pickDataJD:656
京东撞库文件分支"] BR -->|"否"| MAIN["pickData:538 主循环"] subgraph SQLBUILD["SQL 构建 MySQLQueryBuilder"] S1["遍历 condition JSON
conditionMap: 运算符→字段→值"] S1 --> S2["动态字段:t_create_time ≥ -30
→ create_time ≥ '今天-30天'"] S1 --> S3["fileurl 条件 → 不进 WHERE
addFile() 收集文件路径"] S1 --> S4["packageTag 非空 → 拼子查询
contact_number NOT IN
(select phone from people_tag ...)"] S2 --> S5["SELECT selectFields FROM 表
WHERE (c1 AND c2 ...)"] S3 --> S5 S4 --> S5 end QB -.-> SQLBUILD MAIN --> CT["countTotalData
SELECT COUNT(*)"] CT --> BL["布隆过滤器初始化
容量=总量, 误判率按数量级自适应"] BL --> LOOP{"游标分页循环
pageSize=1000
WHERE ... AND id > 上页最大id
ORDER BY id LIMIT 1000"} LOOP --> F1["fetchData:888
敏感字段 AES 解密
(手机号/姓名/身份证/地址)"] F1 --> F2["dataFilter:793
文件匹配:内存 contains 过滤
(运营商文件值必须命中)"] F2 --> F3["deduplicate:864
布隆过滤器按 distinctFields
跨页去重"] F3 --> F4{"超过 limitCount?"} F4 -->|"是"| F5["截断并结束"] F4 -->|"否"| LOOP LOOP -->|"取完"| CSV["writeToCSV:928
清洗引号/换行, 英文逗号→中文
本地文件 = UUID.csv"] CSV --> UP["FTP 上传
dataExtractPath 目录"] UP --> CB["HTTP POST 回调
callbackUrl(原请求全字段回传
+ success/count/remark/ftpPath)"] F5 --> CSV
SQL 细节:所有条件值字符串拼接进 SQL(无参数化);Condition.type(UNION/INTER/DIFF)在本服务中未被读取——多个 Condition 全部 AND 叠加(等价交集),并/差集的真正组合在上游拆分或撞库实现。preCondition 透传但未拼入 SQL。
编排层与 Azkaban 的关系是整个排队与生命周期体系存在的原因:Azkaban 上同一 project 的同一 flow 只能并发跑一个任务,因此编排层用 project 池轮转(acquireProject)争取并发、用 Redis 等待/执行队列缓冲溢出、用 */5s 守护循环做超时兜底与 FIFO 补位。prod 配置 spring.bigdata.url = azkaban.jetmobo.com,project 池 data_extract, data_extract_1~3。
flowchart TD L["1. login()
POST bigdata.url/ 表单
username/password/action=login
→ session.id"] --> P["2. buildParam:479
conditionType=normal → getNormalParam
sql/scala/user_domain → condition 原文透传"] P --> C1["命令文本拼装 buildQuery:544
Equal → name isIn MTIz(Base64)
GreaterThan → age gt MTM=
Between → pt_day between xxx
Like → containsArray / NotLike → notcontainsarray"] C1 --> X["3. extra_param 组装
command + callback + select + header
+ exclude1(固定剔除)/exclude2(自定义排除)
Map → JSON → UTF-8 → Base64 → URLEncoder"] X --> H["4. POST /executor 表单
session.id & ajax=executeFlow
& project={project} & flow={quotaCode}
& flowOverride[extra_param]={base64}"] H --> E2["5. 响应 execid
空或 -1 → failed
成功 → 记入执行 ZSet
project:packageCode:execid"] E2 -.->|"平台异步跑完"| CB2["回调 /api/package/callback
ftpPath 指向大数据侧上传的提数目录"]
人群包生成后按「历史标签」二次排除:下载已有人群包 CSV → 按 phoneType 确定去重列(联系号码/办理号码,或 UUID 模式下的 pu_id/hu_id)→ 每 1000 行一批,把去重列值 AES 加密后与 people_tag 表比对(打过的即已触达)→ 剔除后写新文件 → 回调。本质是按标签历史排除已触达人群。
{dw.site}/api/package/callback(prod 即 http://dc.jetmobo.com);success / count(实际数据量)/ remark(成功为「成功」,失败为异常消息)/ ftpPath;以下问题均来自对当前代码的逐行阅读,建议作为后续治理的输入:
| # | 位置 | 问题 | 影响 |
|---|---|---|---|
| 1 | MySQLQueryBuilder | SQL 字符串拼接,无参数化 | SQL 注入风险(值来自规则配置) |
| 2 | DataPickServiceImpl.dataFilter | 文件匹配用内存 List.contains(O(n) 线性查找) | 大文件场景 O(n²) 性能风险 |
| 3 | execute / distinctByTag | 异步用默认 ForkJoinPool commonPool | 并发任务互相影响,无隔离与背压 |
| 4 | AesUtil | AES 密钥硬编码在源码中 | 密钥泄露即敏感数据(手机号/身份证)暴露 |
| 5 | PickDataManager:294 | datapick 路径的 execid 是本地随机 UUID,非真实执行 id | 无法凭 execid 到平台侧核对任务 |
| 6 | PickDataManager:791 | 超时守护中「杀掉提数调度任务」是 todo | 超时只改本地状态,平台侧任务可能仍在跑 |
| 7 | retryPick:1290 | 重试退避用 Thread.sleep 且在回调 HTTP 线程里睡 1/3/5 分钟 | 回调线程被长时间占用 |
| 8 | retryPick 判断 | 是否重试完全依赖失败 remark 含「重试」字样(下游平台驱动) | 隐式契约,上游失败默认不重试 |
| 9 | datapick-service | pickData 无超时控制;hasHead/preCondition 等字段未消费(表头恒写) | 长任务无兜底;配置字段形同虚设 |
| 10 | datapick-service | 异步失败时本地文件已写入部分数据但回调 success=false 即删除 | 部分结果无法找回 |
| 11 | Condition.type | UNION/INTER/DIFF 在 datapick-service 中未读取,条件全部 AND 叠加 | 多规则「并集/差集」语义在该路径不生效 |