【问题标题】:Spark Analysis Reduce (Twitter Sentiment)Spark 分析减少(Twitter 情绪)
【发布时间】: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 : Longtweet: 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


    【解决方案1】:

    这是一个很好的规范 map-reduce 问题:

    1. 将推文条目映射到表示类别和计数为 1 的元组
    2. 减少类别元组以总结每个类别的数量

    即伪代码:

    + map the RDD you have (id, tweet, pos score...
    - map to a tuple that looks like (category, 1, 1) if the tweet is positive
    - map to a tuple that looks like (category, 1, 0) if the tweet is negative
    
    + reduceByKey where our key is the category using summation
    - we end up with an RDD of tuples in the form you want
    

    这里有一些 scala 代码来完成这个——java 是类似的

    val categoryEntryRDD = tweetsFiltered.map( mappedTuple =>
        if mappedTuple._5 == "positive" {
            (mappedTuple._6, 1, 1)
        } else {
            (mappedTyple._6, 1, 0)
        }
    }
    
    val reducedRDD = categoryEntryRDD.reduceByKey( x, y => (x._1 + y._1, x._2 + y._2) )
    

    此时,reduceRDD 保存的元组看起来像(类别、该类别的推文总数、该类别的正面推文总数)。

    【讨论】:

    • 谢谢,它看起来很棒。但是我现在使用的是 Java,你不能在 JavaRDD 上使用 reduceByKey 方法,有什么想法吗?
    • 我认为它有效,但我有一个问题,如果我有更多类别怎么办?现在它给了我:{TWITTER, 1000, 400},但实际上我也有一个类别 FACEBOOK,例如,您的解决方案只是计算所有内容? :s
    • 是的,最终结果是一个 RDD,其中包含所有唯一类别及其汇总信息的条目。第一个条目可能是 {TWITTER, 1000, 400},第二个条目可能是 {FACEBOOK, 400, 22} 如果有数据的类别为 Facebook。
    • 是的,但 JavaPairRDD 是唯一使用 reduceByKey 方法的 RDD。一个简单的JavaRDD没有reduceByKey方法,只有reduce..但它不起作用。
    【解决方案2】:

    我终于明白了!用Java

    JavaPairRDD<String, Float> categoryPositiveTweets = categoryResult.mapToPair(new PairFunction<Tuple4<Long, String, String, String>, String, Float>() {
            @Override
            public Tuple2<String, Float> call(Tuple4<Long, String, String, String> tuple4) throws Exception {
                if(tuple4._3().equals("positive")){
                    return new Tuple2<String, Float>(tuple4._4(), 1F);
    
                } else {
                    return new Tuple2<String, Float>(tuple4._4(), 0F);
                }
            }
        }).reduceByKey(new Function2<Float, Float, Float>() {
            @Override
            public Float call(Float aFloat, Float aFloat2) throws Exception {
                return aFloat+aFloat2;
            }
        });
    
        JavaPairRDD<String, Float> categoryTotalTweets = categoryResult.mapToPair(new PairFunction<Tuple4<Long, String, String, String>, String, Float>() {
            @Override
            public Tuple2<String, Float> call(Tuple4<Long, String, String, String> tuple4) throws Exception {
                return new Tuple2<String, Float>(tuple4._4(), 1F);
            }
        }).reduceByKey(new Function2<Float, Float, Float>() {
            @Override
            public Float call(Float aFloat, Float aFloat2) throws Exception {
                return aFloat+aFloat2;
            }
        });
    
        JavaPairRDD<String, Tuple2<Float, Float>> joinedCategorizedTweets = categoryTotalTweets.join(categoryPositiveTweets);
    
        JavaRDD<Tuple3<String, Float, Float>> categorizedScoredTweets = joinedCategorizedTweets.map(new Function<Tuple2<String, Tuple2<Float, Float>>, Tuple3<String, Float, Float>>() {
            @Override
            public Tuple3<String, Float, Float> call(Tuple2<String, Tuple2<Float, Float>> tweet) throws Exception {
                return new Tuple3<String, Float, Float>(
                        tweet._1(),
                        tweet._2()._1(),
                        tweet._2()._2());
            }
        });
    

    感谢您的帮助!

    结果:

    (推特, 100, 40) (脸书, 80, 20)

    【讨论】:

      猜你喜欢
      • 2022-11-20
      • 2018-04-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-01-09
      • 2016-10-04
      • 2022-01-10
      相关资源
      最近更新 更多