【发布时间】:2021-02-14 11:30:45
【问题描述】:
我有一堆防火墙数据。 我想:
A) 将每个 IP 每小时的字节数相加,然后
B) 计算该小时内所有 IP 的最小和最大总和
我已经能够在 Kafka 中执行 A,但是,我不知道如何执行 B。我一直在仔细研究文档,感觉快要接近了,但我似乎总是只找到部分解决方案。
我的 firewall_stream 运行良好。
client.create_stream(
table_name='firewall_stream',
columns_type=['src_ip VARCHAR',
'dst_ip VARCHAR',
'src_port INTEGER',
'dst_port INTEGER',
'protocol VARCHAR',
'action VARCHAR',
'timestamp VARCHAR',
'bytes BIGINT',
],
topic='firewall',
value_format='JSON'
)
我创建了物化视图 bytes_sent,滚动窗口为 1 小时,总和(字节)并按 IP 地址分组。这很好用!。
client.ksql('''
CREATE TABLE bytes_sent as
SELECT src_ip, sum(bytes) as bytes_sum
FROM firewall_stream
GROUP BY src_ip
EMIT CHANGES
''')
这就是我卡住的地方。首先,我尝试从 bytes_sent 创建另一个物化视图,该视图通过windowstart 进行了一个 max(bytes_sum) 组,但我收到一个错误,您无法在窗口化的物化视图上进行聚合。
然后我删除了时间窗口(我想我会在第二个物化视图中重新打开它),但是我的“group by”子句没有任何字段。在 Postgres 中,我可以在没有 group by 的情况下做 max,它会在整个表格中计算它,但 Kafka 总是需要那个 group by。现在我不确定要使用什么。
似乎无法与文档中的窗口表格进行连接(尽管我没有尝试过,可能会产生误解)。
我唯一的另一个猜测是从物化视图 bytes_sent 创建另一个流并查看更改日志事件,然后以某种方式将它们转换为给定时间窗口内所有 IP 的最大字节数。
任何有关如何处理此问题的反馈将不胜感激!
【问题讨论】:
标签: apache-kafka apache-kafka-streams ksqldb