【问题标题】:How to count total records read in source using Flink dataset API如何使用 Flink 数据集 API 计算源中读取的总记录数
【发布时间】:2020-08-16 00:06:36
【问题描述】:

我们目前使用 Flink DataSet API 从 FileSystem 读取文件并应用一些批量转换。我们还想获取作业完成后处理的总记录数。 管道就像dataset.map().filter()

count() 函数似乎是一个非并行运算符,它需要从所有数据集中进行额外计算。

是否有任何方法可以计算 map 运算符中已处理的记录并提供像流式传输这样的侧面输出,以便我们可以聚合它们以获得总计数?或者还有其他更好的方法吗?

非常感谢!

【问题讨论】:

  • 看来我需要额外的运算符来进行计数,这意味着我必须对数据集进行两次迭代才能获得原始结果和计数。是否有任何方法可以将计数逻辑集成到地图/平面图运算符中并生成另一个数据集 来进行计数?
  • 我认为这个线程确实回答了这个问题;)这个想法是您需要在一台机器上进行部分计数才能进行最终计数。因此,您需要计算每个键的值,然后计算一台机器上所有键的总计数。

标签: apache-flink


【解决方案1】:

您可能想使用counters。这些计数器允许您为每个任务输出小的统计信息,这些统计信息会在作业完成时自动累积。

【讨论】:

  • 嗨,如何在分离的纱线 flink 作业中获得计数器结果?
  • 计数器通常在您的驱动程序中进行评估(main 调用 execute)。然后,您可以在 execute 完成后以任何您想要的方式发布结果。
  • 似乎只适用于附加模式提交。如果在分离的 mdoe 中向 yarn 提交作业。调用execute()后客户端程序结束,无法得到任何提交结果。
  • 抱歉,没想到客户端在集群之外。一般来说,在分离模式下,你不能在execute之后在客户端执行有意义的事情,除非你等待结果,这违背了分离执行的目的。我会用一些替代方法更新我的答案。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-02-22
  • 1970-01-01
  • 2019-02-01
  • 2012-12-17
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多