求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
相关产品推荐
相关产品推荐

