【问题标题】:Can IF statement work correctly to build spark dataframe?IF 语句可以正常工作以构建 spark 数据框吗?
【发布时间】:2019-05-05 17:58:18
【问题描述】:

我有以下代码,它使用 IF 语句有条件地构建数据框。 这符合我的预期吗?

df = sqlContext.read.option("badRecordsPath", badRecordsPath).json([data_path_1, s3_prefix + "batch_01/2/2019-04-28/15723921/15723921_15.json"])
if "scrape_date" not in df.columns:
    df = df.withColumn("scrape_date", lit(None).cast(StringType()))

【问题讨论】:

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


    【解决方案1】:

    这是你想要做的吗?

    val result = <SOME Dataframe I previously created>
    scala> result.printSchema
    root
     |-- VAR1: string (nullable = true)
     |-- VAR2: double (nullable = true)
     |-- VAR3: string (nullable = true)
     |-- VAR4: string (nullable = true)
    
    scala> result.columns.contains("VAR3")
    res13: Boolean = true
    
    scala> result.columns.contains("VAR9")
    res14: Boolean = false
    
    

    所以“结果”数据框有“VAR1”、“VAR2”等列。 下一行显示它包含“VAR3”(表达式的结果为“true”。但它不包含名为“VAR9”的列(表达式的结果为“false”)。

    上面是 scala,但你应该可以在 Python 中做同样的事情(抱歉,当我回复时,我没有注意到你问的是 python)。

    在执行方面,if语句会在驱动节点本地执行。根据经验,如果某些东西返回 RDD、DataFrame 或 DataSet,它将在 executor(s) 上并行执行。由于 DataFrame.columns 返回一个 Array,所有列列表的处理都将在驱动节点中完成(因为 Array 不是 RDD、DataFrame 也不是 DataSet)。

    还要注意 RDD、DataFrame 和 DataSet 将被“懒惰地”执行。也就是说,Spark 会“累积”生成这些对象的操作。它只会在您执行不生成 RDD、DataFrame 或 DataSet 的操作时执行它们。例如,当您进行表演、计数或收集时。这样做的部分原因是 Spark 可以优化流程的执行。另一个是它只执行生成答案实际需要的操作。

    【讨论】:

    • 感谢您的快速回答。实际上我真正想知道的是: IF 语句是否仅在 spark 驱动程序上执行?或者它会被发射到执行器上运行?
    • df.columns 是一个列表(即不是 RDD、DataFrame 或 DataSet)。因此它在驱动程序上本地执行。在我的示例中,result.columns 是一个 Scala 数组(可能是 Python 中的一个数组)。所以是的,列名的检查将在驱动程序节点中执行。作为一般规则,返回 RDD、DataFrame 或 DataSet 的操作将在集群(executors)中执行
    • 谢谢,这就是我想知道的,现在很清楚了
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-03-14
    • 1970-01-01
    • 2019-07-23
    • 2017-11-25
    • 2016-03-04
    • 2015-09-27
    • 2016-04-14
    相关资源
    最近更新 更多