【问题标题】:Looping a set of queries of in Pyspark running very slow在 Pyspark 中循环运行非常慢的一组查询
【发布时间】:2019-09-13 00:43:29
【问题描述】:

我在 pyspark 中的单个作业中运行 3 个查询。执行第一个查询,然后该查询的结果用于第二个查询,第二个查询的结果用于第三个查询。最后,保存第三个查询。

这项工作需要几个小时。我有一个 15 个实例的 spark 集群和 32 GB 的内存,每个实例有 8 个内核。需要帮助来优化这种情况。

所有这些查询都在循环中运行,如下所示:

对于商店中的 ID: 执行查询 - 1 执行查询 2 执行查询-3

然后保存最终输出。

For 循环迭代 2400 次。

【问题讨论】:

  • 能否包含 Spark 执行计划?例如,通过在 2 次循环后在结果 df 上运行 df.explain()
  • @cylim 我如何在这里发布解释计划?它跨越超过 3-4 页
  • 如果执行计划无法适应问题编辑,您能否显示循环中的 pyspark 代码以及如何检索最终结果?
  • 或者,您可以通过查看此处列出的常见问题来尝试优化您的 PySpark 代码,enigma.com/blog/things-i-wish-id-known-about-spark。我怀疑你的 for 循环中有 pyspark 操作 shuffle 数据,例如 joingroupBy
  • @cylim.. 感谢分享以上链接。

标签: python optimization pyspark apache-spark-sql amazon-emr


【解决方案1】:

迭代使用数据帧很慢

  • 在 for 循环等中迭代使用 Dataframe 很慢,因为它会导致大量查询计划。这个问题详细here

  • 建议的解决方案是转换为 rdd 并再次返回。例如——credit to this post:

    val rdd = df.rdd
    rdd.cache()
    sqlCtx.createDataFrame(rdd. df.schema) 
    
  • 另一种解决方案是使用checkpointing

    From the spark documentation:
    Checkpointing can be used to truncate the logical plan of this DataFrame,
    which is especially useful in iterative algorithms where the plan may grow exponentially. 
    It will be saved to files inside the checkpoint directory set with SparkContext.setCheckpointDir().
    
  • 我不确定这是否已在 spark 3.X 中解决,因为该问题刚刚被标记为已解决,但没有说明采取的任何措施

【讨论】:

    猜你喜欢
    • 2011-08-09
    • 2021-08-19
    • 2018-11-11
    • 1970-01-01
    • 2017-05-12
    • 1970-01-01
    • 2021-05-16
    • 1970-01-01
    • 2016-11-23
    相关资源
    最近更新 更多