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

求NiFi中基于Lastclock属性(大于当前减1分钟)路由FlowFile的Python脚本

NiFi属性路由Python脚本实现

以下是可在NiFi的ExecuteScript处理器中运行的Python脚本,用于根据FlowFile的Lastclock属性(epoch时间戳格式)进行路由:

import time
from org.apache.nifi.processor.io import StreamCallback
from java.io import InputStream, OutputStream

class PropertyRouter(StreamCallback):
    def process(self, inputStream, outputStream):
        # 获取当前FlowFile对象
        flow_file = self.flowFile
        
        # 获取Lastclock属性,若不存在则默认路由到failure
        last_clock_str = flow_file.getAttribute("Lastclock")
        if not last_clock_str:
            self.rel = REL_FAILURE
            return
        
        try:
            # 将属性值转换为整数时间戳
            last_clock = int(last_clock_str)
            # 计算当前时间减去1分钟的时间戳(单位:秒)
            one_minute_ago = int(time.time()) - 60
            
            # 比较时间戳,判断路由方向
            if last_clock > one_minute_ago:
                self.rel = REL_SUCCESS
            else:
                self.rel = REL_FAILURE
        except ValueError:
            # 若Lastclock不是有效整数,路由到failure
            self.rel = REL_FAILURE
        
        # 传递FlowFile到指定关系
        self.flowFile = session.putAttribute(flow_file, "route_result", self.rel.name)
        session.transfer(self.flowFile, self.rel)

# 执行脚本逻辑
flowFile = session.get()
if flowFile is not None:
    flowFile = session.write(flowFile, PropertyRouter())
    session.commit()
else:
    session.commit()

关键逻辑说明

  • 属性获取与校验:先检查Lastclock属性是否存在,若不存在或无法转换为整数时间戳,直接路由至failure队列
  • 时间计算:用time.time()获取当前秒级时间戳,减去60秒得到1分钟前的时间点
  • 路由判断:当Lastclock的时间戳晚于1分钟前的时间时,路由至success,否则走failure

使用注意事项

  • 在NiFi的ExecuteScript处理器中,选择Python作为脚本语言
  • 确保FlowFile的Lastclock属性是有效的秒级epoch时间戳(若为毫秒级,需在转换时除以1000)
  • 处理器需提前配置success和failure两个输出关系

内容的提问来源于stack exchange,提问作者Vyshnavi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 08:08:10