【发布时间】:2017-07-30 11:15:35
【问题描述】:
我们需要像 jaccard 一样在 spark 中的大量数据集上计算距离矩阵。 面临几个问题。请帮我们指路。
问题 1
import info.debatty.java.stringsimilarity.Jaccard;
//sample Data set creation
List<Row> data = Arrays.asList(
RowFactory.create("Hi I heard about Spark", "Hi I Know about Spark"),
RowFactory.create("I wish Java could use case classes","I wish C# could use case classes"),
RowFactory.create("Logistic,regression,models,are,neat","Logistic,regression,models,are,neat"));
StructType schema = new StructType(new StructField[] {new StructField("label", DataTypes.StringType, false,Metadata.empty()),
new StructField("sentence", DataTypes.StringType, false,Metadata.empty()) });
Dataset<Row> sentenceDataFrame = spark.createDataFrame(data, schema);
// Distance matrix object creation
Jaccard jaccard=new Jaccard();
//Working on each of the member element of dataset and applying distance matrix.
Dataset<String> sentenceDataFrame1 =sentenceDataFrame.map(
(MapFunction<Row, String>) row -> "Name: " + jaccard.similarity(row.getString(0),row.getString(1)),Encoders.STRING()
);
sentenceDataFrame1.show();
没有编译时错误。但是出现运行时异常,例如:
org.apache.spark.SparkException: Task not serializable
问题 2
此外,我们需要找到哪对得分最高,我们需要声明一些变量。此外,我们还需要执行其他计算,我们面临很多困难。
即使我尝试在 MapBlock 中声明一个像计数器这样的简单变量,我们也无法捕获增量值。如果我们在 Map 块之外声明,我们会得到很多编译时错误。
int counter=0;
Dataset<String> sentenceDataFrame1 =sentenceDataFrame.map(
(MapFunction<Row, String>) row -> {
System.out.println("Name: " + row.getString(1));
//int counter = 0;
counter++;
System.out.println("Counter: " + counter);
return counter+"";
},Encoders.STRING()
);
请给我们指路。 谢谢。
【问题讨论】:
标签: apache-spark dataset similarity apache-spark-2.0 apache-spark-dataset