【问题标题】:How to generate an attribute according to a condition using executescript in apache nifi (jython)?如何在apache nifi(jython)中使用executescript根据条件生成属性?
【发布时间】:2017-10-16 21:21:30
【问题描述】:

我正在使用 apache nifi 将日志存储在 kafka 中。

我有一些日志行,根据内容我必须将它们发送到主题 kafka 或其他。 我的问题是有很多主题,因此我必须使用很多处理器。

我认为使用'executescript',我可以根据我在日志文本中的条件生成一个动态属性,我可以在publishkafka处理器的'topic name'属性中使用它。

我从代码开始读取流文件的内容,以及写一些条件,但是我不知道如何生成属性。

日志行的一些示例: Jim,18,M,156,俄勒冈,美国等 约翰,55,M,170,爱达荷州,美国等

这是我目前所拥有的:

from org.apache.commons.io import IOUtils
from java.nio.charset import StandardCharsets
from org.apache.nifi.processor.io import StreamCallback

class PyStreamCallback(StreamCallback):
    def __init__(self):
         pass
    def process(self, inputStream, outputStream):
         Log = IOUtils.toString(inputStream, StandardCharsets.UTF_8)
         TextLog = str(Log).split(',')
         name = TextLog[0]
         age = TextLog[1]
         sex = TextLog[2]

         if name == 'John' and age == '30':
             Topic_A = str(TextLog) 
             outputStream.write(bytearray((Topic_A).encode('utf-8')))
         elif name == 'Max' and age == '25':
             Topic_B = str(TextLog)
             outputStream.write(bytearray((Topic_B).encode('utf-8')))
         elif name == 'Smith' and age == '10' or '20':
             Topic_C = str(TextLog)
             outputStream.write(bytearray((Topic_C).encode('utf-8')))

目标是拥有一个 executescript 处理器和一个 kafka 处理器。 拜托,有人可以帮助我。

【问题讨论】:

    标签: apache-nifi


    【解决方案1】:

    简短的回答是:

    flowFile = session.putAttribute(flowFile, 'my-property', 'my-value')
    

    https://funnifi.blogspot.com/2016/02/executescript-processor-hello-world.html

    无需编写任何自定义代码即可完成此操作的典型方法是使用 RouteOnContent...

    您可以添加用户定义的属性,其中名称将成为关系,值是正则表达式。

    例如添加两个属性:

    john = John,30,.*
    max = Max,25,.*
    

    您将从那里将每个关系发送到设置主题名称的 UpdateAttribute 处理器,因此 john 将转到设置 topic = Topic_A 的 UpdateAttribute,而 max 将转到设置 topic = Topic_B 的 UpdateAttribute。

    然后他们都将连接到主题设置为 ${topic} 的单个 PublishKafka。

    【讨论】:

    • 你好布赖恩,因为我有很多kafka主题,为了避免流中有几十个处理器,我尝试用executescript来做。因为我有很多 kafka 主题,如果我使用 routeoncontent 然后更新属性。虽然我只会有一个publishkafka,但我会有很多updateattribute。感谢您的帮助,我会在您发表评论时尝试在putattribute上澄清自己。
    • Brian,如果我理解正确,我需要做的是这样吗? if name == 'John' and age == '30': Topic_A = str(TextLog) outputStream.write(bytearray((Topic_A).encode('utf-8'))) flowFile = session.putAttribute(flowFile, 'my-property', 'Topic_A')
    • 您不想写入 outputStream,因为这样您将覆盖流文件的内容,您想使用输入流回调读取内容,然后适当地设置属性
    • 其实我是打算利用它,删除一些字段。这就是 outputStream.write 的原因。我不清楚的是,如果在同一个脚本中,我可以打开流文件内容,根据某些条件,生成一个属性,从流文件内容中删除一些字段,最后让流文件输出已经根据相应的属性修改它的内容。
    • 不能在回调中添加属性,看看这个例子,先写内容,写完回调后,再添加属性funnifi.blogspot.com/2016/02/…
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-06-12
    • 1970-01-01
    • 1970-01-01
    • 2016-08-30
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多