如何基于HTTP响应使用Groovy脚本停止StreamSets管道
StreamSets 通过Groovy评估器触发自定义事件停止管道实现方案
链路配置逻辑
管道按以下顺序连接即可同时满足「记录正常写入文件」「读到最终记录自动停管道」两个需求:
- HTTP Client源 → Groovy评估器常规数据输出端口 → 文件写入目标阶段
- Groovy评估器事件输出端口 → Pipeline Finisher(管道终结器,StreamSets原生自带执行器,无需额外开发)
前置确认项
先明确HTTP接口返回JSON中标记拉取结束的特征,常见的判断维度有三种,选和你业务匹配的即可:
- 单条记录携带结束标识:比如
isFinal: true字段 - 分页接口无下一页:比如
nextPageUrl字段为null或空字符串 - 批次返回空数据:接口拉取到的记录列表长度为0
Groovy评估器实现脚本
评估器默认开启常规数据输出,记得在评估器配置页手动开启事件输出端口,脚本直接复制后修改判断逻辑即可:
import com.streamsets.pipeline.api.Event import com.streamsets.pipeline.api.Record // 遍历当前批次所有输入记录 for (Record record : records) { // 所有记录优先写入常规输出,保证下游文件阶段能收到全量数据 output.write(record) // -------------------------- // 替换成你实际的最终记录判断逻辑 boolean hitFinalFlag = record.get('/isFinal').getValueAsBoolean() // 分页场景判断逻辑参考: // boolean hitFinalFlag = (record.get('/nextPage').getValueAsString() == null || record.get('/nextPage').getValueAsString().isBlank()) // 空批次判断逻辑参考(把判断逻辑移到for循环外,写if(records.size() == 0)再触发事件即可) // -------------------------- if (hitFinalFlag) { // 创建自定义停止事件 Event stopEvent = eventCreator.createEvent("custom-stop", 1) // 可按需给事件附加元信息,非必须 stopEvent.set("/stopTrigger", "HTTP源返回最终记录,正常终止") stopEvent.set("/triggerRecordId", record.get('/id').getValueAsString()) // 事件输出到事件端口,触发后续Pipeline Finisher执行 eventOutput.write(stopEvent) } }
避坑说明
- 不要在Groovy脚本里直接调用Runtime或者StreamSets硬停止接口,这种方式会直接杀进程,当前批次未写完的记录会丢失,用事件触发Pipeline Finisher的方式会等所有在途记录写入完成后再正常停管道。
- 测试阶段先把Groovy评估器的事件输出接到Trash(丢弃阶段),先验证全量记录能正常写入文件、结束标识判断逻辑准确,再切换接到Pipeline Finisher,避免逻辑错误导致管道提前终止。
- 如果HTTP源返回的结束标识是在响应头而非JSON体里,要先在HTTP Client源配置里把需要的响应头提取成记录属性,脚本里通过
record.getHeader().getAttribute('你提取的头字段名')取值判断即可。
内容的提问来源于stack exchange,提问作者AJ20
相关产品推荐
相关产品推荐

