【发布时间】: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