Spark Session read mulitple files instead of using pattern

Viewed 3409

I'm trying to read couple of CSV files using SparkSession from a folder on HDFS ( i.e I don't want to read all the files in the folder )

I get the following error while running (code at the end):

Path does not exist:
file:/home/cloudera/works/JavaKafkaSparkStream/input/input_2.csv,
/home/cloudera/works/JavaKafkaSparkStream/input/input_1.csv

I don't want to use the pattern while reading, like /home/temp/*.csv, reason being in future I have logic to pick only one or two files in the folder out of 100 CSV files

Please advise

    SparkSession sparkSession = SparkSession
            .builder()
            .appName(SparkCSVProcessors.class.getName())
            .master(master).getOrCreate();
    SparkContext context = sparkSession.sparkContext();
    context.setLogLevel("ERROR");

    Set<String> fileSet = Files.list(Paths.get("/home/cloudera/works/JavaKafkaSparkStream/input/"))
            .filter(name -> name.toString().endsWith(".csv"))
            .map(name -> name.toString())
            .collect(Collectors.toSet());

    SQLContext sqlCtx = sparkSession.sqlContext();

    Dataset<Row> rawDataset = sparkSession.read()
            .option("inferSchema", "true")
            .option("header", "true")
            .format("com.databricks.spark.csv")
            .option("delimiter", ",")
            //.load(String.join(" , ", fileSet));
            .load("/home/cloudera/works/JavaKafkaSparkStream/input/input_2.csv, " +
                    "/home/cloudera/works/JavaKafkaSparkStream/input/input_1.csv");

UPDATE

I can iterate the files and do an union as below. Please recommend if there is a better way ...

    Dataset<Row> unifiedDataset = null;

    for (String fileName : fileSet) {
        Dataset<Row> tempDataset = sparkSession.read()
                .option("inferSchema", "true")
                .option("header", "true")
                .format("csv")
                .option("delimiter", ",")
                .load(fileName);
        if (unifiedDataset != null) {
            unifiedDataset= unifiedDataset.unionAll(tempDataset);
        } else {
            unifiedDataset = tempDataset;
        }
    }
3 Answers
Related