【问题标题】:Load csv data with partition in spark 2.0在 spark 2.0 中加载带有分区的 csv 数据
【发布时间】:2016-08-26 22:54:06
【问题描述】:

在 Spark 2.0 中,我有以下方法将数据加载到数据集中

public Dataset<AccountingData> GetDataFrameFromTextFile()
{     // The schema is encoded in a string
    String schemaString = "id firstname lastname accountNo";

    // Generate the schema based on the string of schema
    List<StructField> fields = new ArrayList<>();
    for (String fieldName : schemaString.split("\t")) {
        StructField field = DataTypes.createStructField(fieldName, DataTypes.StringType, true);
        fields.add(field);
    }
    StructType schema = DataTypes.createStructType(fields);

    return  sparksession.read().schema(schema)
            .option("mode", "DROPMALFORMED")
            .option("sep", "|")
            .option("ignoreLeadingWhiteSpace", true)
            .option("ignoreTrailingWhiteSpace ", true)
            .csv("D:\\HadoopDirectory\Employee.txt").as(Encoders.bean(Employee.class));
}

在我的驱动程序代码中,对数据集调用 Map 操作

    Dataset<Employee> rowDataset = ad.GetDataFrameFromTextFile();

    Dataset<String> map = rowDataset.map(new MapFunction<Employee, String>() {
        @Override
        public String call(Employee emp) throws Exception {
            return  TraverseRuleByADRow(emp);
        }
    },Encoders.STRING());

当我在笔记本电脑上以 8 个内核在 spark 本地模式下运行驱动程序时,我看到 8 个分区分割了输入文件。请问是否有办法将文件加载到 8 个以上的分区中,比如 100 个或1000 个分区?

如果源数据通过 jdbc 来自 sql server 表,我知道这是可以实现的。

sparksession.read().format("jdbc").option("url", urlCandi).option("dbtable", tableName).option("partitionColumn", partitionColumn).option("lowerBound", String.valueOf(lowerBound))
                .option("upperBound", String.valueOf(upperBound))
                .option("numPartitions", String.valueOf(numberOfPartitions))
                .load().as(Encoders.bean(Employee.class));

谢谢

【问题讨论】:

    标签: java apache-spark


    【解决方案1】:

    使用 Dataset 中的 repartition() 方法。根据Scaladoc,在读取时没有设置分区数的选项

    【讨论】:

    • 在我的代码中,我使用了 coalesce(1,true) 输出到单个文件。我仍然看到机器上可用的内核数的分区数。如果我使用了 coalesce(1000,true),我会在输出目录中看到 1000 个小文件,但将文件拆分为 1000 个分区并没有帮助。 sourceDataSet.toJavaRDD().coalesce(1,true).saveAsTextFile(OUTPUT_DIRECTORY);
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-07-15
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多