【发布时间】:2017-09-19 21:50:06
【问题描述】:
我使用的是 Spark 2.1.0。对于以下代码,它读取文本文件并将内容转换为 DataFrame,然后输入 Word2Vector 模型:
SparkSession spark = SparkSession.builder().appName("word2vector").getOrCreate();
JavaRDD<String> lines = spark.sparkContext().textFile("input.txt", 10).toJavaRDD();
JavaRDD<List<String>> lists = lines.map(new Function<String, List<String>>(){
public List<String> call(String line){
List<String> list = Arrays.asList(line.split(" "));
return list;
}
});
JavaRDD<Row> rows = lists.map(new Function<List<String>, Row>() {
public Row call(List<String> list) {
return RowFactory.create(list);
}
});
StructType schema = new StructType(new StructField[] {
new StructField("text", new ArrayType(DataTypes.StringType, true), false, Metadata.empty())
});
Dataset<Row> input = spark.createDataFrame(rows, schema);
input.show(3);
Word2Vec word2Vec = new Word2Vec().setInputCol("text").setOutputCol("result").setVectorSize(100).setMinCount(0);
Word2VecModel model = word2Vec.fit(input);
Dataset<Row> result = model.transform(input);
抛出异常
java.lang.RuntimeException:编码时出错:java.util.Arrays$ArrayList 不是有效的外部类型 数组架构
这发生在 input.show(3) 行,因此 createDataFrame() 导致异常,因为 Arrays.asList() 返回一个 Arrays$ArrayList 此处不支持。但是 Spark 官方文档有以下代码:
List<Row> data = Arrays.asList(
RowFactory.create(Arrays.asList("Hi I heard about Spark".split(" "))),
RowFactory.create(Arrays.asList("I wish Java could use case classes".split(" "))),
RowFactory.create(Arrays.asList("Logistic regression models are neat".split(" ")))
);
StructType schema = new StructType(new StructField[]{
new StructField("text", new ArrayType(DataTypes.StringType, true), false, Metadata.empty())
});
Dataset<Row> documentDF = spark.createDataFrame(data, schema);
效果很好。如果 Arrays$ArrayList 不受支持,那么这段代码是如何工作的?不同之处在于我将JavaRDD<Row> 转换为DataFrame,但官方文档将List<Row> 转换为DataFrame。我相信 Spark Java API 有一个重载方法 createDataFrame(),它接受 JavaRDD<Row> 并根据提供的模式将其转换为 DataFrame。我很困惑为什么它不起作用。任何人都可以帮忙吗?
【问题讨论】:
标签: java apache-spark