【发布时间】:2023-03-14 04:40:02
【问题描述】:
也许这是一个愚蠢的问题,但我不得不问。
我在 Nifi 中有一个 Collect_data 处理器,它将消息流式传输到另一个进程,该进程使用 python 脚本来解析并创建 json 文件。问题是我不知道 python 脚本中函数的输入是什么。如何将这些消息(16 位数字)从 Collect_data 处理器传递到下一个处理器包含 python 脚本。有什么好的,基本的例子吗?
我已经在网上找了一些例子,但不是很明白。
import datetime
import hashlib
from urlparse import urlparse, parse_qs
import sys
from urlparse import urlparse, parse_qs
from datetime import *
import json
import java.io
from org.apache.commons.io import IOUtils
from java.nio.charset import StandardCharsets
from org.apache.nifi.processor.io import StreamCallback
from time import time
def parse_zap(inputStream, outputStream):
data = inputStream
buf = (hashlib.sha256(bytearray.fromhex(data)).hexdigest())
buf = int(buf, 16)
buf_check = str(buf)
if buf_check[17] == 2:
pass
datetime_now = datetime.now()
log_date = datetime_now.isoformat()
try:
mac = buf_check[7:14].upper()
ams_id = buf_check[8:]
action = buf_check[3:4]
time_a = int(time())
dict_test = {
"user": {
"guruq" : 'false'
},
"device" : {
"type" : "siolbox",
"mac": mac
},
"event" : {
"origin" : "iptv",
"timestamp": time_a,
"type": "zap",
"product-type" : "tv-channel",
"channel": {
"id" : 'channel_id',
"ams-id": ams_id
},
"content": {
"action": action
}
}
}
return dict_test
except Exception as e:
print('%s nod PARSE 500 \"%s\"' % (log_date, e))
感谢我没看错,但现在我无法创建输出。 提前致谢。
【问题讨论】:
标签: python arguments apache-nifi