【问题标题】:process values of records [duplicate]记录的过程值[重复]
【发布时间】:2020-06-10 21:14:36
【问题描述】:

我是 Spark 的新手,我找不到足够的信息来理解 Spark 中的某些内容。我正在尝试在 scala 中编写伪代码(例如这些示例http://spark.apache.org/examples.html

给出了一个包含数据的文件。每行都有一些数据:编号、课程名称、学分和分数。

123 Programming_1 10 75
123 History       5  80

我正在尝试计算每个学生(人数)的平均值。平均是每门课程学分的总和*学生所学的分数 除以学生修读的每门课程学分的总和。忽略任何具有 mark==NULL 的行。假设我有一个函数 parseData(line),它用字符串创建一行以记录 4 个成员:数字、课程名称、学分、标记。

我到现在为止的尝试

data=spark.textFile(“hdfs://…”)
line=data.filter(mark=> mark != null)
line= line.map(line => parseData(line))
data = parallelize(List(line))  
groupkey= data.groupByKey()  
               ((a,b,c)=>(a, sum(mul(b,c))/ sum(b))

但我不知道如何读取具体值并使用它们来计算每个学生的平均值。可以用数组吗?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql rdd


    【解决方案1】:

    过滤并获取数据框后,您可以使用以下内容:

    df.withColumn("product",col("credits")*col("marks"))
    .groupBy(col("student"))
    .agg(sum("credits").as("sumCredits"),sum("product").as("sumProduct"))
    .withColumn("average",col("sumProduct")/col("sumCredits"))
    

    希望这会有所帮助!

    【讨论】:

    • 我不明白。它是什么语言?你能像这里的这些例子一样吗? spark.apache.org/examples.html
    • 这是用scala写的
    • spark 中有两个抽象级别 - RDD(低级 api)和数据帧和数据集(高级 api)。以上代码用于数据框方法
    • 我明白了。此代码对于数据成员是否完整并忽略任何具有标记==NULL 的行? df.withColumn 和 agg 有什么作用?
    • 不 - 这不会过滤空值,您可以使用 df.filter / df.where 过滤掉空值。在获得学分和分数的产品之前执行此操作。
    猜你喜欢
    • 1970-01-01
    • 2020-02-18
    • 1970-01-01
    • 1970-01-01
    • 2015-11-20
    • 1970-01-01
    • 2013-01-26
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多