【发布时间】: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,然后将其存储?
【问题讨论】: