【问题标题】:Spark SQL .withColumn() vs Column expressionsSpark SQL .withColumn() 与列表达式
【发布时间】:2019-11-23 01:26:10
【问题描述】:

我想知道在 pyspark 中使用 中间步骤/列 时是否有任何性能/可扩展性差异:

  1. 使用 .withColumn() 例如:
    df = df.withColumn('bar', df.foo + 1)
    df = df.withColumn('baz', df.bar + 2)

然后拨打df.select('baz').collect()

  1. 将 Spark 列声明为 Python 变量:
    bar = df.foo + 1
    baz = bar + 2

然后调用 df.select(baz.alias('baz')).collect()

问题:如果需要很多中间步骤/列,例如bar,这两个选项在空间/时间复杂度上会有所不同吗?

【问题讨论】:

  • 我看到我的帖子被删除了。事后看来,这很可能是正确的,因为我引用的示例有 foldLeft。除非主持人缺乏沟通。你的例子

标签: apache-spark pyspark apache-spark-sql pyspark-sql


【解决方案1】:

我看到我原来的帖子被删除了。事后看来,这很可能是正确的,除非缺乏沟通。该示例使用的是 foldLeft,它不是您的融合数据管道的用例。

为了回答您的问题,Catalyst 融合数据管道操作意味着物理计划显示的任何一种方式都不存在性能问题:

df = spark.createDataFrame([(x,x) for x in range(7)], ['foo', 'bar',])
df = df.withColumn('bar', df.foo + 1) 
df = df.withColumn('baz', df.bar + 2)
df.select('baz').explain(extended=True)

== Physical Plan ==
*(1) Project [(foo#276L + 3) AS baz#283L]
+- *(1) Scan ExistingRDD[foo#276L,bar#277L]  

同样:

df = spark.createDataFrame([(x,x) for x in range(7)], ['foo', 'bar',])
bar = df.foo + 1 
baz = bar + 2
df.select(baz.alias('baz')).explain(extended=True)

== Physical Plan ==
*(1) Project [(foo#288L + 3) AS baz#292L]
+- *(1) Scan ExistingRDD[foo#288L,bar#289L]

它们看起来和我很相似……注意 +3 的优化。

此外,我提请您注意使用 foldLeft 和 .withColumn https://manuzhang.github.io/2018/07/11/spark-catalyst-cost.html

【讨论】:

    猜你喜欢
    • 2017-12-03
    • 2020-05-04
    • 1970-01-01
    • 2018-12-26
    • 1970-01-01
    • 1970-01-01
    • 2021-11-09
    • 2019-05-30
    • 1970-01-01
    相关资源
    最近更新 更多