【发布时间】:2018-09-20 01:21:07
【问题描述】:
我有以下数据框:df
在某些时候,我需要根据时间戳(毫秒)过滤掉项目。 但是,保存过滤了多少记录对我来说很重要(如果太多,我想让工作失败) 天真地我能做到:
======Lots of calculations on df ======
val df_filtered = df.filter($"ts" >= startDay && $"ts" <= endDay)
val filtered_count = df.count - df_filtered.count
但感觉完全是矫枉过正,因为 SPARK 将执行整个执行树 3 次(过滤和 2 个计数)。 Hadoop MapReduce 中的这项任务非常简单,因为我可以为过滤的每一行维护计数器。 有没有更有效的方法,只能找到累加器却无法连接过滤器。
建议的方法是在过滤器之前缓存 df,但是由于 DF 大小,我更喜欢这个选项作为最后的手段。
【问题讨论】:
-
与计数相比,这种方法如何更快?有没有办法精确计算?它仍然没有处理计数本身之前的所有工作,这将在没有缓存的情况下执行所有计算的 3 倍
-
df.except(df_filtered).count怎么样? -
累加器使用示例可参考here。但请注意,在转换中使用累加器并不是 100% 准确,因为在失败的情况下可以多次执行转换。但我想在你的情况下,这件事可以忽略。另一个缺点是当你发现错误太多时不能立即停止处理,因为在过滤器操作中无法获取累加器值。所以你需要在所有处理后失败。
-
@VladislavVarslavans 您能否给出简单的示例或链接,说明如何将累加器与数据框过滤器一起使用,因为我看不到它。如果我过滤这些值,我该如何为它们存储累加器
标签: scala apache-spark dataframe bigdata