【问题标题】:Best way to read TSV file using Apache Spark in java在 java 中使用 Apache Spark 读取 TSV 文件的最佳方法
【发布时间】:2016-12-09 15:07:08
【问题描述】:

我有一个 TSV 文件,其中第一行是标题。我想从这个文件创建一个 JavaPairRDD。目前,我正在使用以下代码:

TsvParser tsvParser = new TsvParser(new TsvParserSettings());
List<String[]> allRows;
List<String> headerRow;
try (BufferedReader reader = new BufferedReader(new FileReader(myFile))) {
        allRows = tsvParser.parseAll((reader));
        //Removes the header row
        headerRow = Arrays.asList(allRows.remove(0));
    }
JavaPairRDD<String, MyObject> myObjectRDD = javaSparkContext
            .parallelize(allRows)
            .mapToPair(row -> new Tuple2<>(row[0], myObjectFromArray(row)));

我想知道是否有办法让 javaSparkContext 直接读取和处理文件,而不是将操作分成两部分。

编辑:这不是 How do I convert csv file to rdd 的副本,因为我正在寻找 Java 中的答案,而不是 Scala。

【问题讨论】:

标签: java csv apache-spark


【解决方案1】:

使用https://github.com/databricks/spark-csv

import org.apache.spark.sql.SQLContext

SQLContext sqlContext = new SQLContext(sc);
DataFrame df = sqlContext.read()
    .format("com.databricks.spark.csv")
    .option("inferSchema", "true")
    .option("header", "true")
    .option("delimiter","\t")
    .load("cars.csv");

df.select("year", "model").write()
    .format("com.databricks.spark.csv")
    .option("header", "true")
    .save("newcars.csv");

【讨论】:

  • 我查看了它,但我不喜欢将 Dataframe 转换为自定义对象的 JavaRDD 的简单方法。您必须手动解析每一行。
  • 我很好奇为什么在唯一对象要求的实例中你会想要一个数据帧上的 RDD。 @alexgbelov
  • 我的其余逻辑使用 RDD;我只需要一种尽可能干净地读取数据的方法。
  • spark 2x 实现该功能。你不再需要 spark-csv
【解决方案2】:

尝试下面的代码来读取 CSV 文件并创建 JavaPairRDD。

public class SparkCSVReader {

public static void main(String[] args) {

    SparkConf conf = new SparkConf().setAppName("CSV Reader");
    JavaSparkContext sc = new JavaSparkContext(conf);
    JavaRDD<String> allRows = sc.textFile("c:\\temp\\test.csv");//read csv file
    String header = allRows.first();//take out header
    JavaRDD<String> filteredRows = allRows.filter(row -> !row.equals(header));//filter header
    JavaPairRDD<String, MyCSVFile> filteredRowsPairRDD = filteredRows.mapToPair(parseCSVFile);//create pair
    filteredRowsPairRDD.foreach(data -> {
        System.out.println(data._1() + " ### " + data._2().toString());// print row and object
    });
    sc.stop();
    sc.close();
}

private static PairFunction<String, String, MyCSVFile> parseCSVFile = (row) -> {
    String[] fields = row.split(",");
    return new Tuple2<String, MyCSVFile>(row, new MyCSVFile(fields[0], fields[1], fields[2]));
};

}

您还可以使用 Databricks spark-csv (https://github.com/databricks/spark-csv)。 spark-csv 也包含在 Spark 2.0.0 中。

【讨论】:

  • 谢谢。我查看了 Databricks spark-csv,但我不喜欢它,因为您必须手动解析每一行,反正我已经在这样做了。
【解决方案3】:

Apache Spark 2.x 具有内置的 csv 阅读器,因此您不必使用 https://github.com/databricks/spark-csv

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;

/**
 *
 * @author cpu11453local
 */
public class Main {
    public static void main(String[] args) {


        SparkSession spark = SparkSession.builder()
                .master("local")
                .appName("meowingful")
                .getOrCreate();

        Dataset<Row> df = spark.read()
                    .option("header", "true")
                    .option("delimiter","\t")
                    .csv("hdfs://127.0.0.1:9000/data/meow_data.csv");

        df.show();
    }
}

还有maven文件pom.xml

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <groupId>com.meow.meowingful</groupId>
    <artifactId>meowingful</artifactId>
    <version>1.0-SNAPSHOT</version>
    <packaging>jar</packaging>
    <properties>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
        <maven.compiler.source>1.8</maven.compiler.source>
        <maven.compiler.target>1.8</maven.compiler.target>
    </properties>

    <dependencies>
        <!-- https://mvnrepository.com/artifact/org.apache.spark/spark-core_2.11 -->
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.11</artifactId>
            <version>2.2.0</version>
        </dependency>


        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-sql_2.11</artifactId>
            <version>2.2.0</version>
        </dependency>
    </dependencies>

</project>

【讨论】:

    【解决方案4】:

    我是uniVocity-parsers 的作者,在 spark 方面帮不上什么忙,但我相信这样的东西对你有用:

    parserSettings.setHeaderExtractionEnabled(true); //captures the header row
    
    parserSettings.setProcessor(new AbstractRowProcessor(){
            @Override
            public void rowProcessed(String[] row, ParsingContext context) {
                String[] headers = context.headers() //not sure if you need them
                JavaPairRDD<String, MyObject> myObjectRDD = javaSparkContext
                        .mapToPair(row -> new Tuple2<>(row[0], myObjectFromArray(row)));
                //process your stuff.
            }
        });
    

    如果要并行处理每一行,可以换一个ConcurrentRowProcessor

    parserSettings.setProcessor(new ConcurrentRowProcessor(new AbstractRowProcessor(){
            @Override
            public void rowProcessed(String[] row, ParsingContext context) {
                String[] headers = context.headers() //not sure if you need them
                JavaPairRDD<String, MyObject> myObjectRDD = javaSparkContext
                        .mapToPair(row -> new Tuple2<>(row[0], myObjectFromArray(row)));
                //process your stuff.
            }
        }, 1000)); //1000 rows loaded in memory.
    

    然后调用解析:

    new TsvParser(parserSettings).parse(myFile);
    

    希望这会有所帮助!

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-02-09
      • 2019-04-04
      • 1970-01-01
      • 2017-03-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-09-17
      相关资源
      最近更新 更多