【问题标题】:Java Spark Streaming with Cassandra使用 Cassandra 进行 Java Spark 流式处理
【发布时间】:2017-06-26 19:54:51
【问题描述】:

我正在尝试使用 Cassandra 进行 java spark 流式传输。我对 Scala 做了同样的事情,但我不知道如何在 Java 中进行。网络没有给我任何关于 java spark 流和 Cassandra 的例子。

有人可以告诉我如何在 java 中使用以下 Scala 代码:

import org.apache.spark.streaming.dstream.ConstantInputDStream

val ssc = new StreamingContext(conf, Seconds(10))

val cassandraRDD = ssc.cassandraTable("mykeyspace", "users").select("fname", "lname").where("lname = ?", "yu")

val dstream = new ConstantInputDStream(ssc, cassandraRDD)

dstream.foreachRDD{ rdd => 
    // any action will trigger the underlying cassandra query, using collect to have a simple output
    println(rdd.collect.mkString("\n")) 
}
ssc.start()
ssc.awaitTermination()

感谢任何帮助。谢谢

【问题讨论】:

    标签: java scala cassandra spark-streaming rdd


    【解决方案1】:

    在您的 foreachRDD 转换中,您可以按照 cassandra 表格式转换数据。

    JavaRDD<TestBean> cassandraRDD = testRDD
                    .flatMap(new FlatMapFunction<Tuple2<String, List<Map<String, Object>>>, TestBean>() {
    
                        private static final long serialVersionUID = 1L;
    
                        @Override
                        public Iterable<TestBean> call(Tuple2<String, List<Map<String, Object>>> tuple) throws Exception {
    
                            return rawData;
                        }
                    });
    
                javaFunctions(jsonRDD).writerBuilder(CASSANDRA_KEYSPACE,CASSANDRA_TABLE, mapToRow(TestBean.class)).saveToCassandra();
    

    【讨论】:

    • 这些是我的导入语句 import static com.datastax.spark.connector.japi.CassandraJavaUtil.javaFunctions;导入静态 com.datastax.spark.connector.japi.CassandraJavaUtil.mapToRow;
    猜你喜欢
    • 1970-01-01
    • 2015-09-05
    • 2019-12-10
    • 2020-08-11
    • 1970-01-01
    • 2023-04-03
    • 2016-08-14
    • 2018-01-11
    • 2016-09-28
    相关资源
    最近更新 更多