【发布时间】:2019-05-22 19:06:54
【问题描述】:
我正在尝试使用 MongoSpark 和 rdd (JavaMongoRdd) 在 java 中执行 mapReduce。所以目前,我可以在我的 Rdd 中检索我的 mongo 文档,但我不知道之后如何继续。事实上,我的文档中有一个字段,它是一个日期,我想使用这个日期中的年份来执行我的 mapReduce,但我没有找到任何关于如何执行此操作的信息。所以我在这里问你是否有一些文档、教程,甚至是如何进行的示例。
这里的代码,我正在尝试使用带有 Mongo 文档和年份的 pairRdd 来计算每年的文档数量,但我不知道这是否是我必须继续的方式
public String count() {
JavaSparkContext jsc = new JavaSparkContext(sparkSession.sparkContext());
JavaMongoRDD<Document> rdd = MongoSpark.load(jsc);
logger.info("test 1 :" + rdd.count());
logger.info("test 2 :" + rdd.first().toJson());
/*JavaMongoRDD<Document> newRdd = rdd.withPipeline(
Collections.singletonList(
Document.parse("{ $match: { _id : { $gt : ObjectId(\"5c9e180cdba48525f0df30b9\") } } }")
)
);*/
//logger.info("test 2.5 :" +newRdd.first());
JavaPairRDD<String, Document> pairRdd = rdd
.mapToPair((document) -> new Tuple2(document.getString("date").split(".")[1], document));
logger.info("test 3 :" + pairRdd.first());
//logger.info("test 2 :" + rdd.first().toJson());
//ar
//logger.info("test spark");
return "test";
}
我的 MongoDb 文档是这样的
"_id" : ObjectId("5c9e180ddba48525f0df30cb"),
"title" : "Redevance: une perte de compétitivité pour l’hydraulique suisse",
"description" : [
"Le Parlement a bouclé, durant cette session de printemps, la révision de la loi sur les forces hydrauliques. La solution adoptée aboutit au statu quo sur le plan de la redevance hydraulique. Le taux maximal de cette taxe reste ainsi fixé à 110 francs par kilowatt théorique, jusqu'à fin 2024. Les..."
],
"date" : "dimanche, 24. mars 2019"
【问题讨论】:
-
提供一些代码。
-
我添加了我已经完成的一点点代码
-
是获取日期还是使用映射和归约函数的问题?
-
问题在于使用地图和归约函数,当我谈到日期时只是为了解释我在做什么