【发布时间】:2016-03-15 03:52:31
【问题描述】:
我有一个关于 apache Spark 和 Java 的问题
我正在制作一个从 Twitter (Twitter4J) 流式传输数据的应用程序。我也在制作一个分析数据的应用程序。一个带有 JSON 推文的 txt 文件。
流媒体应用: 输出 tweet.txt: 示例:一行Json:
{"id":674534622903054336,"user":"twitter","tweet":"a tweet from twitter #twitter.","date":"2015-12-09T11:22:41CET"}
AnalyzerApp:
SparkConf conf = new SparkConf().setMaster("local[2]").setAppName("TwitterAnalyzerBigData");
final JavaSparkContext sc = new JavaSparkContext(conf);
JavaRDD<String> jsonFile = sc.textFile("whateverpath/tweets.txt");
JavaPairRDD<Long, String> tweetsFiltered = jsonFile.mapToPair(new TwitterFilterFunction());
tweetsFiltered 是一个 JavaPairRDD:tweet ID : Long 和 tweet: String
现在我正在使用一些地图函数来获得这样的结果:
(1,a tweet from twitter #twitter.,0.0,0.055555556,negative, TWITTER)
(这是随机测试数据)
- 1 是 ID
- 来自推特#twitter 的推文:推文
- 0.0:正分
- 0.0566:负分
- 负面:类别情绪(正面或负面)
- TWITTER:推文类别(基于主题标签的类别)
问题:我怎样才能减少这个RDD,所以我得到这样的结果:
TWITTER, 1, 0
- TWITTER:推文的类别
- 1 : TWITTER CATEGORY 的推文总数
- 0:TWITTER CATEGORY 的正面推文数量
在 James 的回答之后,我用 Java 制作了 reduceByKey。
JavaRDD<Tuple3<String, Float, Float>> categoryEntryRDD = categoryResult.map(new Function<Tuple4<Long, String, String, String>, Tuple3<String, Float, Float>>() {
@Override
public Tuple3<String, Float, Float> call(Tuple4<Long, String, String, String> tuple4) throws Exception {
if(tuple4._3().equals("positive")){
return new Tuple3<String, Float, Float>(tuple4._4(), 1F, 1F);
} else {
return new Tuple3<String, Float, Float>(tuple4._4(), 1F, 0F);
}
}
});
Tuple3<String, Float, Float> reducedRDD = categoryEntryRDD.reduce(new Function2<Tuple3<String, Float, Float>, Tuple3<String, Float, Float>, Tuple3<String, Float, Float>>() {
@Override
public Tuple3<String, Float, Float> call(Tuple3<String, Float, Float> tuple31, Tuple3<String, Float, Float> tuple32) throws Exception {
System.out.println(tuple31.toString());
return new Tuple3<String, Float, Float>(tuple31._1(), tuple31._2()+tuple32._2(), tuple31._3()+tuple32._3());
}
});
但是reduce方法和reduceByKey不一样,怎么解决呢?
我的输出: {推特, 1000, 400} 但我也有一个类别:拥有 1000 条推文的 FACEBOOK。
【问题讨论】:
标签: java twitter apache-spark mapreduce sentiment-analysis