【问题标题】:Creating an RDD after retrieving data from cassandra DB从 cassandra DB 检索数据后创建 RDD
【发布时间】:2015-10-21 10:59:17
【问题描述】:

我正在为我的项目使用 cassandra 和 spark,现在我写这个是为了从数据库中检索数据:

 results = session.execute("SELECT * FROM foo.test");

 ArrayList<String> supportList = new ArrayList<String>();
 for (Row row : results) {
            supportList.add(row.getString("firstColumn") + "," + row.getString("secondColumn")));
        }
        JavaRDD<String> input = sparkContext.parallelize(supportList);
        JavaPairRDD<String, Double> tuple = input.mapToPair(new PairFunction<String, String, Double>() {
            public Tuple2<String, Double> call(String x) {
                String[] parts = x.split(",");
                return new Tuple2(parts[0],String.valueOf(new Random().nextInt(30) + 1));
            }

它工作,但我想知道是否有一个漂亮的方法来编写上面的代码,我想要实现的是:

  • 在 scala 中,我可以通过这种方式简单地检索和填充 RDD:

    val dataRDD = sc.cassandraTable[TableColumnNames]("keySpace", "table")

  • 如何在不使用支持列表或其他“讨厌”的东西的情况下用 Java 编写相同的东西。

更新

JavaRDD<String> cassandraRowsRDD = javaFunctions(javaSparkContext).cassandraTable("keyspace", "table")
                .map(new Function<CassandraRow, String>() {
                    @Override
                    public String call(CassandraRow cassandraRow) throws Exception {
                        return cassandraRow.toString();
                    }
                });

我在这一行 -> public String call(CassandraRow cassandraRow) 这个例外:

Exception in thread "main" org.apache.spark.SparkException: Task not serializable
    at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:166)
    at org.apache.spark.util.ClosureCleaner$.clean(ClosureCleaner.scala:158)
    at org.apache.spark.SparkContext.clean(SparkContext.scala:1623)
    at org.apache.spark.rdd.RDD.map(RDD.scala:286)
    at org.apache.spark.api.java.JavaRDDLike$class.map(JavaRDDLike.scala:89)
    at org.apache.spark.api.java.AbstractJavaRDDLike.map(JavaRDDLike.scala:46)
    at org.sparkexamples.cassandraExample.main.KMeans.executeQuery(KMeans.java:271)
    at org.sparkexamples.cassandraExample.main.KMeans.main(KMeans.java:67)
Caused by: java.io.NotSerializableException: org.sparkexamples.cassandraExample.main.KMeans
Serialization stack:
    - object not serializable (class: org.sparkexamples.cassandraExample.main.KMeans, value: org.sparkexamples.cassandraExample.main.KMeans@3015db78)
    - field (class: org.sparkexamples.cassandraExample.main.KMeans$2, name: this$0, type: class org.sparkexamples.cassandraExample.main.KMeans)
    - object (class org.sparkexamples.cassandraExample.main.KMeans$2, org.sparkexamples.cassandraExample.main.KMeans$2@5dbf5634)
    - field (class: org.apache.spark.api.java.JavaPairRDD$$anonfun$toScalaFunction$1, name: fun$1, type: interface org.apache.spark.api.java.function.Function)
    - object (class org.apache.spark.api.java.JavaPairRDD$$anonfun$toScalaFunction$1, <function1>)
    at org.apache.spark.serializer.SerializationDebugger$.improveException(SerializationDebugger.scala:38)
    at org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:47)
    at org.apache.spark.serializer.JavaSerializerInstance.serialize(JavaSerializer.scala:80)
    at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:164)
    ... 7 more

提前致谢。

【问题讨论】:

  • 如果你想完全像在你的 scala 示例中那样做,为什么不使用 Java API? datastax.com/dev/blog/accessing-cassandra-from-spark-in-java
  • @ccheneson 我看到了这些 api。你能看到问题更新吗,我收到一个错误。
  • 我没有看到任何问题更新。请复制/粘贴您在帖子中遇到的错误
  • 你的 scala 示例有效吗?
  • 你能显示正在使用的导入吗?

标签: java cassandra apache-spark rdd


【解决方案1】:

看看答案:RDD not serializable Cassandra/Spark connector java API

问题可能是您显示的代码块周围的类不是可序列化的。

【讨论】:

  • 在我发布的链接中,该类确实实现了Serializable。所以可能是public class JavaDemo implements Serializable {
【解决方案2】:

我遇到了同样的问题。我在一个单独的类中实现了 spark 接口函数,并将其提供给地图功能。它在那之后起作用了。

样本

public a 实现函数 {....}

在地图中使用了这个

.....地图(新a())

它得到了纠正。关于匿名类的火花反序列化的一些问题。

【讨论】:

    猜你喜欢
    • 2018-02-23
    • 2014-01-22
    • 1970-01-01
    • 1970-01-01
    • 2021-12-10
    • 2019-07-14
    • 1970-01-01
    • 2017-05-21
    • 2015-09-12
    相关资源
    最近更新 更多