【问题标题】:Trying to use map on a Spark DataFrame尝试在 Spark DataFrame 上使用地图
【发布时间】:2017-07-22 12:18:56
【问题描述】:

我最近开始尝试使用 Spark 和 Java。我最初使用RDD 浏览了著名的WordCountexample,一切都按预期进行。现在我正在尝试实现我自己的示例,但使用的是 DataFrames 而不是 RDD。

所以我正在从文件中读取数据集

DataFrame df = sqlContext.read()
        .format("com.databricks.spark.csv")
        .option("inferSchema", "true")
        .option("delimiter", ";")
        .option("header", "true")
        .load(inputFilePath);

然后我尝试选择一个特定的列并像这样对每一行应用一个简单的转换

df = df.select("start")
        .map(text -> text + "asd");

但是编译发现第二行有问题,我不完全理解(开始列推断为string类型)。

在接口 scala.Function1 中找到多个非覆盖抽象方法

为什么我的 lambda 函数被视为 Scala 函数,错误消息的实际含义是什么?

【问题讨论】:

    标签: java apache-spark java-8 apache-spark-sql spark-dataframe


    【解决方案1】:

    如果您在数据帧上使用select函数,您会得到一个数据帧。然后在Rowdatatype 上应用一个函数,而不是行的值。之后您应该首先获取该值,因此您应该执行以下操作:

    df.select("start").map(el->el.getString(0)+"asd")

    但是你会得到一个 RDD 作为返回值而不是一个 DF

    【讨论】:

    • 你还需要在map之前有.javaRDD()。另外据我了解,没有明确的方法可以将函数应用于 DataFrame 的行并取回 DataFrame 对吗?
    • 或许你可以试试:df.select("start").forEach(el->el.getString(0)+"asd")
    【解决方案2】:

    我使用 concat 来实现这个

    df.withColumn( concat(col('start'), lit('asd'))
    

    由于您两次映射相同的文本,我不确定您是否还希望替换字符串的第一部分?但如果你是,我会这样做:

    df.withColumn('start', concat(
                          when(col('start') == 'text', lit('new'))
                          .otherwise(col('start))
                         , lit('asd')
                         )
    
    

    此解决方案在使用大数据时可扩展,因为它连接两列而不是迭代值。

    【讨论】:

    • 能否请您告诉我如何在 java stackoverflow.com/questions/63668096/… 中处理这个用例 ...
    • 我在工作atm,但可以稍后看看,但先看看这个答案stackoverflow.com/questions/63537324/…
    • 非常感谢您的快速回复,当然,请指导我修复这个用例....如果第一列是,我需要用另一列的值替换一列的值null ...此处发布的链接显示值替换为另一个。你能帮我吗
    • 我无法在一个集合循环中执行多个 withColumns ...这会导致两个数据集...
    猜你喜欢
    • 2018-04-17
    • 2020-08-10
    • 2023-03-21
    • 1970-01-01
    • 2020-08-09
    • 1970-01-01
    • 1970-01-01
    • 2020-10-29
    • 2019-06-09
    相关资源
    最近更新 更多