【问题标题】:Cassandra Spark ConnectorCassandra 火花连接器
【发布时间】:2016-12-09 03:11:26
【问题描述】:

我的 cassandra CF 有 date 和 id 作为 partition Key 。 查询时我只知道日期,所以我循环了 id 的范围。

我的问题围绕连接器如何执行以下代码。

SparkDriver 代码看起来像 -

SparkConf conf = new SparkConf().setAppName("DemoApp")
.conf.setMaster("local[*]")
.set("spark.cassandra.connection.host", "10.*.*.*")
.set("spark.cassandra.connection.port", "*");

JavaSparkContext sc = new JavaSparkContext(conf);
SparkContextJavaFunctions javaFunctions = CassandraJavaUtil.javaFunctions(sc);

String date = "23012017";

for(String id : idlist) {

JavaRDD<CassandraRow> cassandraRowsRDD = 

javaFunctions.cassandraTable("datakeyspace", "sample2")
            .where("date = ?",date)
            .where("id = ? ", id)
            .select("data");

 cassandraRowsRDDList.add(cassandraRowsRDD);
}

List<CassandraRow> collectAllRows = new ArrayList<CassandraRow>();
        for(JavaRDD<CassandraRow> rdd : cassandraRowsRDDList){
            //do transformations

            collectAllRows.addAll(rdd.collect());
    }

1) 首先我想问一下我是否循环遍历 idlist,比如 idlist 有 1000 个元素,这些元素可能会不断增加,这会有效吗?每个选择查询如何分布在集群中?尤其是如何维护 Cassandra DB 连接?

2) 在我的驱动程序中循环后,我将所有行放在 List 中,然后对每一行应用转换并过滤掉重复项。这是否也会通过集群上的 spark 进行分发,还是会发生在驾驶员一侧。

请帮忙!

【问题讨论】:

    标签: apache-spark cassandra-2.0 datastax-enterprise datastax-java-driver spark-cassandra-connector


    【解决方案1】:

    spark cassandra 连接器提供了一种更好的方法。 您可以创建一个 (date,id) 的 rdd,然后在列 date 和 id 上调用 joinWithCassandraTable 函数。连接器做得很巧妙,所有数据都将仅由工作人员获取,而且无需随机播放,每个工作人员将仅获取其拥有的日期和 ID 的数据。

    【讨论】:

      猜你喜欢
      • 2016-07-30
      • 2018-08-18
      • 1970-01-01
      • 2020-10-04
      • 2019-06-01
      • 1970-01-01
      • 2017-07-15
      • 2016-09-06
      • 2016-05-01
      相关资源
      最近更新 更多