You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 23:48:22