【问题标题】:How to apply map function on in spark java RDD Operations如何在 Spark Java RDD 操作中应用地图功能
【发布时间】:2017-04-29 11:07:21
【问题描述】:

我的 CSV 文件:

YEAR,UTILITY_ID,UTILITY_NAME,OWNERSHIP,STATE_CODE,AMR_METERING_RESIDENTIAL,AMR_METERING_COMMERCIAL,AMR_METERING_INDUSTRIAL,AMR_METERING_TRANS,AMR_METERING_TOTAL,AMI_METERING_RESIDENTIAL,AMI_METERING_COMMERCIAL,AMI_METERING_INDUSTRIAL,AMI_METERING_TRANS,AMI_METERING_TOTAL,ENERGY_SERVED_RESIDENTIAL,ENERGY_SERVED_COMMERCIAL,ENERGY_SERVED_INDUSTRIAL,ENERGY_SERVED_TRANS,ENERGY_SERVED_TOTAL
2011,34,City of Abbeville - (SC),M,SC,880,14,,,894,,,,,,,,,,
2011,84,A & N Electric Coop,C,MD,135,25,,,160,,,,,,,,,,
2011,84,A & N Electric Coop,C,VA,31893,2107,0,,34000,,,,,,,,,,
2011,97,Adams Electric Coop,C,IL,8334,190,,,8524,,,,,0,,,,,0
2011,108,Adams-Columbia Electric Coop,C,WI,33524,1788,709,,36021,,,,,,,,,,
2011,118,Adams Rural Electric Coop, Inc,C,OH,7457,20,,,7477,,,,,,,,,,
2011,122,Village of Arcade,M,NY,3560,498,100,,4158,,,,,,,,,,
2011,155,Agralite Electric Coop,C,MN,4383,227,315,,4925,,,,,,,,,,

下面是读取 CSV 文件的 Spark 代码:

import org.apache.spark.api.java.JavaSparkContext;

public class RddCsv 
{
    public static void main(String[] args) 
    {
    SparkConf conf = new SparkConf().setAppName("CSV Reader").setMaster("local");
    JavaSparkContext sc = new JavaSparkContext(conf);
    JavaRDD<String> allRows = sc.textFile("file:///home/kumar/Desktop/Eletricaldata/file8_2011.csv");//read csv file
    System.out.println(allRows.take(5)); 
   }
}

我是 sparkJava 新手, 如何从该 CsvDataset 中选择 Perticuler 字段值以及如何执行聚合操作,以及如何使用给定数据集的转换和操作。以及如何选择特定的字段值

【问题讨论】:

标签: java csv apache-spark


【解决方案1】:
public static void main(String[] args)
{
    SparkConf conf = new SparkConf().setAppName("CSV Reader").setMaster("local");
    JavaSparkContext sc = new JavaSparkContext(conf);
    JavaRDD<String> allRows = sc.textFile("file:///home/abhishek/Desktop/file8_2011.csv");
    System.out.println(allRows.take(5));
    List<String> headers= Arrays.asList(allRows.take(1).get(0).split(","));
    String field="YEAR";
    //Skip Header
    JavaRDD<String>dataWithoutHeaders=allRows.filter(x -> !(x.split(",")[headers.indexOf(field)]).equals(field));
    //Take one field as integer
    JavaRDD<Integer> years=dataWithoutHeaders.map(x -> Integer.valueOf(x.split(",")[headers.indexOf(field)]));
    //Aggregate operation getTotal aggregate() arguments are initial value for a partition,aggregating function for a partition
    //and aggregating function for results from different partition
    int total=years.aggregate(0,RddCsv::sum,RddCsv::sum);
    for (Integer i:years.collect()){
        System.out.println("year :: "+i);
    }
    System.out.println(total);
}

private static int sum(int a,int b){
    return a+b;
}

这是一个基本程序。您应该阅读 spark 的 java api 以获取详细信息。

【讨论】:

  • 不工作,compitime 错误正在这一行 .aggregate(0,RddCsv::sum,RddCsv::sum);
  • 我刚刚运行过了。
  • 输出:16/12/14 16:26:31 INFO TaskSetManager:在本地主机(1/1)16/12/14 16 上的 25 毫秒内完成阶段 3.0(TID 3)中的任务 0.0: 26:31 INFO TaskSchedulerImpl:从池 16/12/14 中删除了 TaskSet 3.0,其任务已全部完成 16:26:31 INFO DAGScheduler:作业 3 完成:在 RDDCsv.java:32 收集,耗时 0.049729 年 :: 2011年 :: 2011 年 :: 2011 年 :: 2011 年 :: 2011 年 :: 2011 年 :: 2011 年 :: 2011 16088 16/12/14 16:26:31 信息 SparkContext:从关闭挂钩调用 stop()
  • 请检查您的配置、文件名等。或在评论中粘贴您的错误。
  • 这一行: int total=years.aggregate(0,RddCsv::sum,RddCsv::sum);List headers= Arrays.asList(allRows.take(1).get( 0).split(","));
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-05-07
  • 2019-02-11
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多