使用Control M调度Nifi process group时如何确认数据已全部入库完成
NiFi作为流处理框架,原生设计上没有批处理场景下的“执行完成”状态标识,你需要通过自定义标记或者API统计指标来判断批次任务是否执行完成。
方案1:流程内加计数标记(准确性最高,推荐生产使用)
- 在写入数据库的处理器(如PutDatabaseRecord)后新增计数逻辑,用UpdateAttribute分别统计源文件总记录数和成功写入数据库的记录数
- 当两个计数值相等时,用PutFile生成独立的完成标记文件,写入本地指定目录,文件名可携带批次时间、流程组ID等唯一标识
- 调用启动流程组的shell脚本中新增轮询逻辑:
- 检测指定目录下是否存在对应批次的完成标记文件
- 额外查询目标数据库的对应批次写入条数,和源文件总条数做二次校验
- 校验全部通过后再调用API停止流程组,返回成功状态给Control M即可
方案2:仅通过NiFi API轮询(无需修改现有流程,适合临时场景)
- 轮询NiFi流程组状态API,路径为
/nifi-api/process-groups/{你的流程组ID}/status - 同时满足以下两个条件时判定为批次执行完成:
- 接口返回的
processGroupStatus.aggregateSnapshot.queuedCount值为0,即当前流程组内没有排队待处理的流文件 - 连续3次轮询(间隔建议10秒)该值都保持为0,且此前已有成功流出的流文件记录
- 接口返回的
- 注意该方案存在边界风险:如果流文件处理报错进入失败队列/死信队列,也会出现排队数为0的情况,建议搭配数据库侧条数校验使用
核心操作命令示例
启动流程组
curl -X PUT -H "Content-Type: application/json" http://<NiFi主机地址>:<端口>/nifi-api/process-groups/<流程组ID>/run-status -d '{"state": "RUNNING"}'
查询流程组排队流文件数
# 需提前安装jq工具用于解析json curl -s http://<NiFi主机地址>:<端口>/nifi-api/process-groups/<流程组ID>/status | jq -r '.processGroupStatus.aggregateSnapshot.queuedCount'
轮询判断完成的示例脚本
pg_id="替换为你的流程组ID" nifi_addr="http://<NiFi主机地址>:<端口>" # 启动流程组 curl -X PUT -H "Content-Type: application/json" ${nifi_addr}/nifi-api/process-groups/${pg_id}/run-status -d '{"state": "RUNNING"}' # 轮询等待队列清空 while true; do queued_cnt=$(curl -s ${nifi_addr}/nifi-api/process-groups/${pg_id}/status | jq -r '.processGroupStatus.aggregateSnapshot.queuedCount') if [ "${queued_cnt}" = "0" ]; then # 连续三次检测避免临时为空的误判 sleep 10 queued_cnt2=$(curl -s ${nifi_addr}/nifi-api/process-groups/${pg_id}/status | jq -r '.processGroupStatus.aggregateSnapshot.queuedCount') sleep 10 queued_cnt3=$(curl -s ${nifi_addr}/nifi-api/process-groups/${pg_id}/status | jq -r '.processGroupStatus.aggregateSnapshot.queuedCount') if [ "${queued_cnt2}" = "0" ] && [ "${queued_cnt3}" = "0" ]; then break fi fi sleep 30 done # 队列清空后停止流程组 curl -X PUT -H "Content-Type: application/json" ${nifi_addr}/nifi-api/process-groups/${pg_id}/run-status -d '{"state": "STOPPED"}' # 此处可添加自定义数据库条数校验逻辑,校验通过后返回0给Control M
内容的提问来源于stack exchange,提问作者Chakkara
相关产品推荐
相关产品推荐

