【问题标题】:In PySpark groupBy, how do I calculate execution time by group?在 PySpark groupBy 中,如何按组计算执行时间?
【发布时间】:2021-07-22 03:53:02
【问题描述】:

我在一个大学项目中使用 PySpark,我有大型数据框,我使用 groupBy 应用了 PandasUDF。基本上调用是这样的:

df.groupBy(col).apply(pandasUDF)

我在我的 Spark 配置中使用 10 个内核 (SparkConf().setMaster('local[10]'))。

目标是能够报告每个组运行我的代码所花费的时间。我想要每组完成的时间,以便我可以取平均值。我也对计算标准差感兴趣。

我现在正在使用我知道将分成 10 组的已清理数据进行测试,并且我让 UDF 使用 time.time() 打印运行时间。但是,如果我要使用更多组,这是不可能的(对于上下文,我的所有数据都将被分成 3000 个左右的组)。有没有办法衡量每个组的执行时间?

【问题讨论】:

  • 要计算并报告每次调用 UDF 的执行时间,我认为您需要像当前一样在 UDF 中执行此操作。如果您想要总执行时间,您可以将计算出的执行时间添加到 Spark 累加器,然后在应用程序结束时打印它。 spark.apache.org/docs/3.1.1/…
  • 感谢您的回复。不幸的是,我对总时间不感兴趣。我也可以在我的笔记本上轻松做到这一点。我想要每组完成的时间,以便我可以取平均值。也许我可以尝试保存在变量中,但我不确定这是否会发生,我的 udf 必须返回其他内容。
  • 您不必返回花费的时间。只需打印它,然后检查容器日志(如果是 Yarn)。
  • 如果不想将执行时间打印到标准输出,那么从 Pandas UDF 中将其作为额外列返回是否可以满足您的需求?
  • 是的,将相同的数字添加到组中的所有行是我的想法以及我可能会做的事情,除非我找到一种方法来在 udf 运行时填充本地列表然后保存归档。

标签: apache-spark pyspark apache-spark-sql user-defined-functions


【解决方案1】:

如果不想将执行时间打印到标准输出,您可以将其作为 Pandas UDF 的额外列返回,例如

@pandas_udf("my_col long, execution_time long", PandasUDFType.GROUPED_MAP)
def my_pandas_udf(pdf):
    start = datetime.now()
    # Some business logic
    return pdf.assign(execution_time=datetime.now() - start)

或者,要计算驱动程序应用程序的平均执行时间,您可以使用两个Accumulators 在 UDF 中累积执行时间和 UDF 调用次数。例如

udf_count = sc.accumulator(0)
total_udf_execution_time = sc.accumulator(0)

@pandas_udf("my_col long", PandasUDFType.GROUPED_MAP)
def my_pandas_udf(pdf):
    start = datetime.now()
    # Some business logic
    udf_count.add(1)
    total_udf_execution_time.add(datetime.now() - start)
    return pdf

# Some Spark action to run business logic

mean_udf_execution_time = total_udf_execution_time.value / udf_count.value

【讨论】:

  • 感谢累加器的建议,它看起来非常有用。如果我还想获得运行时间的标准偏差,我想我将不得不采用第一种方法(添加新的 col)对吗?
  • 是的,我想是的。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2022-11-17
  • 2017-12-27
  • 2018-03-07
  • 2013-04-30
  • 1970-01-01
  • 1970-01-01
  • 2016-10-01
相关资源
最近更新 更多