【发布时间】:2020-08-16 00:06:36
【问题描述】:
我们目前使用 Flink DataSet API 从 FileSystem 读取文件并应用一些批量转换。我们还想获取作业完成后处理的总记录数。
管道就像dataset.map().filter()
count() 函数似乎是一个非并行运算符,它需要从所有数据集中进行额外计算。
是否有任何方法可以计算 map 运算符中已处理的记录并提供像流式传输这样的侧面输出,以便我们可以聚合它们以获得总计数?或者还有其他更好的方法吗?
非常感谢!
【问题讨论】:
-
看来我需要额外的运算符来进行计数,这意味着我必须对数据集进行两次迭代才能获得原始结果和计数。是否有任何方法可以将计数逻辑集成到地图/平面图运算符中并生成另一个数据集
来进行计数? -
我认为这个线程确实回答了这个问题;)这个想法是您需要在一台机器上进行部分计数才能进行最终计数。因此,您需要计算每个键的值,然后计算一台机器上所有键的总计数。
标签: apache-flink