【问题标题】:spark scala : Convert DataFrame OR Dataset to single comma separated stringspark scala:将 DataFrame OR Dataset 转换为单个逗号分隔的字符串
【发布时间】:2023-04-05 14:26:01
【问题描述】:

下面是 spark scala 代码,它将打印一列 DataSet[Row]:

import org.apache.spark.sql.{Dataset, Row, SparkSession}
val spark: SparkSession = SparkSession.builder()
        .appName("Spark DataValidation")
        .config("SPARK_MAJOR_VERSION", "2").enableHiveSupport()
        .getOrCreate()

val kafkaPath:String="hdfs:///landing/APPLICATION/*"
val targetPath:String="hdfs://datacompare/3"
val pk:String = "APPLICATION_ID" 
val pkValues = spark
        .read
        .json(kafkaPath)
        .select("message.data.*")
        .select(pk)
        .distinct() 
pkValues.show()

关于代码的输出:

+--------------+
|APPLICATION_ID|
+--------------+
|           388|
|           447|
|           346|
|           861|
|           361|
|           557|
|           482|
|           518|
|           432|
|           422|
|           533|
|           733|
|           472|
|           457|
|           387|
|           394|
|           786|
|           458|
+--------------+

问题:

如何将此数据框转换为逗号分隔的字符串变量?

预期输出:

val   data:String= "388,447,346,861,361,557,482,518,432,422,533,733,472,457,387,394,786,458"

请建议如何将 DataFrame[Row] 或 Dataset 转换为一个 String 。

【问题讨论】:

    标签: java scala apache-spark spark-dataframe


    【解决方案1】:

    我认为这不是一个好主意,因为 dataFrame 是一个分布式对象并且可能非常庞大。 Collect 会将所有数据带到驱动程序中,请谨慎执行此类操作。

    以下是您可以使用 dataFrame 执行的操作(两个选项):

    df.select("APPLICATION_ID").rdd.map(r => r(0)).collect.mkString(",")
    df.select("APPLICATION_ID").collect.mkString(",")
    

    只有 3 行的测试数据帧的结果:

    String = 388,447,346
    

    编辑:使用 DataSet 你可以直接做:

    ds.collect.mkString(",")
    

    【讨论】:

    • 非常感谢您的快速回复。 ds.collect.mkString(",") 工作并在最终值中找到 [] 所以替换方法将其删除
    • 即使使用 DataFrame 也无需使用 RDD API(因为他似乎使用 spark 2),您可以省略 .rdd 甚至使用 df.select("APPLICATION_ID").as[String].collect.mkString(",")
    • 你不能只是.mkString(",") 转义单元格可能需要。
    • 这适用于非流数据集。如何使用无法使用 collect 的流式源实现类似的功能。更多详情,stackoverflow.com/questions/62746964/…
    【解决方案2】:

    使用collect_list:

    import org.apache.spark.sql.functions._
    val data = pkValues.select(collect_list(col(pk))) // collect to one row
        .as[Array[Long]] // set encoder, so you will have strongly-typed Dataset
        .take(1)(0) // get the first row - result will be Array[Long]
        .mkString(",") // and join all values
    

    但是,执行收集或获取所有行是一个非常糟糕的主意。相反,您可能希望将 pkValues 保存在 .write?或者将其作为其他函数的参数,以保持分布式计算

    编辑:刚刚注意到,@SCouto 在我之后发布了其他答案。收集也是正确的,使用 collect_list 函数你有一个优势 - 如果你愿意,你可以轻松地进行分组,即将键分组为偶数和奇数。取决于您喜欢哪种解决方案,使用 collect 更简单,或者更长一行,但更强大

    【讨论】:

    • 因为他需要不同的值,您也可以使用collect_set,然后删除distinct。而不是take(1)(0),你可以使用first
    猜你喜欢
    • 2020-09-02
    • 2012-02-01
    • 1970-01-01
    • 1970-01-01
    • 2019-11-18
    • 2018-08-16
    • 2013-09-30
    • 1970-01-01
    • 2021-12-13
    相关资源
    最近更新 更多