【问题标题】:Store countByKey result into Cassandra将 countByKey 结果存储到 Cassandra
【发布时间】:2015-05-11 09:54:31
【问题描述】:

我想计算每个用户在任何一天(来自 Cassandra 表)的 IndicatePresence 消息数量,然后将其存储在单独的 Cassandra 表中以驱动一些仪表板页面。我设法让 'countByKey' 工作,但现在无法弄清楚如何将 Spark-Cassandra 'saveToCassandra' 方法与 Map 一起使用(它只需要 RDD)。

    JavaSparkContext sc = new JavaSparkContext(conf);
    CassandraJavaRDD<CassandraRow> indicatePresenceTable = javaFunctions(sc).cassandraTable("mykeyspace", "indicatepresence");
    JavaPairRDD<UserDate, CassandraRow> keyedByUserDate = indicatePresenceTable.keyBy(new Function<CassandraRow, UserDate>() {
        private static final long serialVersionUID = 1L;
        @Override
        public UserDate call(CassandraRow cassandraIndicatePresenceRow) throws Exception {
            SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
            return new UserDate(cassandraIndicatePresenceRow.getString("userid"), sdf.format(cassandraIndicatePresenceRow.getDate("date")));
        }
    });

    Map<UserDate, Object> countByKey = keyedByUserDate.countByKey();

    writerBuilder("analytics", "countbykey", ???).saveToCassandra();

有没有办法直接在 writerBuilder 中使用 Map?或者我应该编写自己的自定义reducer,它返回一个RDD,但本质上与countByKey方法做同样的事情?或者,我是否应该将 Map 中的每个条目转换为新的 POJO(例如 UserDateCount,带有用户、日期和计数)并使用“并行化”将列表转换为 RDD,然后将其存储?

【问题讨论】:

    标签: cassandra apache-spark


    【解决方案1】:

    最好的办法是永远不要将结果返回给驱动程序(通过使用 countByKey)。而是使用 reduceByKey 以 (key, count) 的形式获取另一个 RDD。将该 RDD 映射到表格的行格式,然后在其上调用 saveToCassandra

    这种方法最重要的优势是我们从不将数据序列化回驱动程序应用程序。所有信息都保存在集群上,并从它们直接保存到 C*,而不是通过驱动程序应用程序的瓶颈运行。

    示例(非常类似于 Map Reduce Word Count):

    1. 将每个元素映射到 (key, 1)
    2. 调用reduceByKey改变(key, 1) -> (key, count)
    3. 将每个元素映射到可写到 C* (key,count)-> WritableObject 的东西
    4. 调用保存到 C*

    在 Scala 中,这类似于

    keyedByUserDate
      .map(_.1, 1)                               // Take the Key portion of the tuple and replace the value portion with 1
      .reduceByKey( _ + _ )                      // Combine the value portions for all elements which share a key
      .map{ case (key, value) => your C* format} // Change the Tuple2 to something that matches your C* table
      .saveToCassandra(ks,tab)                   // Save to Cassandra
    

    在 Java 中它有点复杂(为 K 和 V 插入你的类型)

    .mapToPair(new PairFunction<Tuple2<K,V>,K,Long>>, Tuple2<K, Long>(){
        @Override
        public Tuple2<K, Long> call(Tuple2<K, V> input) throws Exception {
          return new Tuple2(input._1(),1)
        }
    }.reduceByKey(new Function2(Long,Long,Long)(){
        @Override
        public Long call(Long value1, Long value2) throws Exception {
          return value1 + value2
        }
    }.map(new Function1(Tuple2<K, Long>, OutputTableClass)(){  
        @Override
        public OutputTableClass call(Tuple2<K,Long> input) throws Exception {
        //Do some work here
        return new OutputTableClass(col1,col2,col3 ... colN)
       }
    }.saveToCassandra(ks,tab, mapToRow(OutputTableClass.class))
    

    【讨论】:

    • 啊,好吧,这就是我的想法——所以基本上,如果我找到“countByKey”的源代码并将其放入我自己的“reduceByKey”自定义函数中,但将其保留为 RDD 而不是让它调用'mapAsSerializableJavaMap',这样可以吗?
    • 就是这个想法,reduceByKey 应该只是一种 (_ + _) 的 lambda
    • 抱歉,我对此很陌生,我对 reduce 感到困惑 - 对于 RDD,reduce 接受 并返回 V。这是否意味着要进行计数,我的 V 类必须在自身内部包含一个计数字段(大概在运行 reduce 之前它将被硬编码为 1)?在我的 V 是 CassandraRow 的情况下,大概我将不得不创建一个像 CassandraRowWithCount 这样的包装器并在我的 RDD 中使用它?还是有更好的方法?
    • 我在问题中添加更多细节
    • 啊,返回的地图 (key, 1) 有点我想不通 - 很好的答案,感谢您的帮助!
    猜你喜欢
    • 1970-01-01
    • 2019-06-16
    • 2021-04-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-04-21
    相关资源
    最近更新 更多