【问题标题】:Saving JavaList to Cassandra table using spark context使用火花上下文将 JavaList 保存到 Cassandra 表
【发布时间】:2016-10-03 00:10:09
【问题描述】:

嗨,我是 spark 和 scala 的新手,我在将数据保存到 cassandra 时遇到了一些问题,下面是我的场景

1) 我从我的 java 类到 scala 类获取用户定义对象的列表(比如包含名字、姓氏等的用户对象),到这里为止,我可以访问用户对象并能够打印它的内容

2) 现在我想使用 spark 上下文将该 usersList 保存到 cassandra 表中,我已经经历了许多示例,但在每个地方我都看到使用我们的 caseClass 创建 Seq 和硬编码值,然后保存到 cassandra,我已经尝试过,并且对我来说工作正常,如下所示

import scala.collection.JavaConversions._
import org.apache.spark.SparkConf
import org.apache.spark.SparkContext

import com.datastax.spark.connector._
import java.util.ArrayList

object SparkCassandra extends App {
    val conf = new SparkConf()
        .setMaster("local[*]")
        .setAppName("SparkCassandra")
        //set Cassandra host address as your local address
        .set("spark.cassandra.connection.host", "127.0.0.1")
    val sc = new SparkContext(conf)
     val usersList = Test.getUsers
     usersList.foreach(x => print(x.getFirstName))
    val collection = sc.parallelize(Seq(userTable("testName1"), userTable("testName1")))
    collection.saveToCassandra("demo", "user", SomeColumns("name"))
    sc.stop()
}

case class userTable(name: String)

但这里我的要求是使用来自我的 usersList 的动态值,而不是硬编码值,或任何其他方式来实现这一点。

【问题讨论】:

  • 有多少用户?这些值存储在哪里?
  • 将有多达 20k 用户,actullay 我从其他一些 javaClass 获得该列表并需要存储在 cassandra 表中
  • 只要你在并行化,它应该可以工作。如何从“usersList”创建一个包含“userTable”所有案例类对象的 Seq 并并行化并保存?
  • 如果您可以发布错误以查看究竟出了什么问题,那就很容易了。
  • @Sreekar 我没有收到任何错误,我正在寻找将列表数据插入到 cassandra 表中的各种方法

标签: java scala apache-spark cassandra


【解决方案1】:

如果您创建CassandraRow 对象的RDD,则可以直接保存结果,而无需指定列或案例类。此外,CassandraRow 具有非常方便的fromMap 函数,因此您可以将行定义为Map 对象,对其进行转换并保存结果。

例子:

val myData = sc.parallelize(
  Seq(
    Map("name" -> "spiffman", "address" -> "127.0.0.1"),
    Map("name" -> "Shabarinath", "address" -> "127.0.0.1")
  )
)

val cassandraRowData = myData.map(rowMap => CassandraRow.fromMap(rowMap))

cassandraRowData.saveToCassandra("keyspace", "table")

【讨论】:

  • 感谢回复,这里我的要求是不要使用硬编码值代替“spiffman”和“shabarinath”我需要使用列表对象值
  • 列表对象的类型是什么?你能把它转换成地图吗?
  • 简单 pojo 用户对象列表,其中包含名字、getter 和 setter,我想存储该列表
  • 我的意思是Test.getUsers的输出类型是什么?如果是Seq[String]的用户名,你可以做sc.parallelize(Test.getUsers.map(user => CassandraRow.fromMap(Map("name" -> user)).saveToCassandra("keyspace", "table")for example。
  • getUsers 的输出是 List (java.util.list)
【解决方案2】:

最后我得到了满足我的要求的解决方案,并且工作正常,如下所示:

我的 Scala 代码:

import scala.collection.JavaConversions.asScalaBuffer
import scala.reflect.runtime.universe
import org.apache.spark.SparkConf
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD
import com.datastax.spark.connector.SomeColumns
import com.datastax.spark.connector.toNamedColumnRef
import com.datastax.spark.connector.toRDDFunctions

object JavaListInsert {
  def randomStores(sc: SparkContext, users: List[User]): RDD[(String, String, String)] = {
       sc.parallelize(users).map { x => 
       val fistName = x.getFirstName
       val lastName = x.getLastName
       val city = x.getCity
       (fistName, lastName, city)
    }
  }

  def main(args: Array[String]): Unit = {
    val conf = new SparkConf().setAppName("cassandraInsert")
    val sc = new SparkContext(conf)
    val usersList = Test.getUsers.toList
    randomStores(sc, usersList).
      saveToCassandra("test", "stores", SomeColumns("first_name", "last_name", "city"))
    sc.stop
  }
}

Java Pojo 对象:

    import java.io.Serializable;
    public class User implements Serializable{
        private static final long serialVersionUID = -187292417543564400L;
        private String firstName;
        private String lastName;
        private String city;

        public String getFirstName() {
            return firstName;
        }

        public void setFirstName(String firstName) {
            this.firstName = firstName;
        }

        public String getLastName() {
            return lastName;
        }

        public void setLastName(String lastName) {
            this.lastName = lastName;
        }

        public String getCity() {
            return city;
        }

        public void setCity(String city) {
            this.city = city;
        }
}

返回用户列表的 Java 类:

import java.util.ArrayList;
import java.util.List;


public class Test {
    public static List<User> getUsers() {
        ArrayList<User> usersList = new ArrayList<User>();
        for(int i=1;i<=100;i++) {
            User user = new User();
            user.setFirstName("firstName_+"+i);
            user.setLastName("lastName_+"+i);
            user.setCity("city_+"+i);
            usersList.add(user);
        }
        return usersList;
    }
}

【讨论】:

    猜你喜欢
    • 2019-02-27
    • 2018-02-11
    • 2017-01-01
    • 1970-01-01
    • 2016-02-07
    • 2017-07-05
    • 1970-01-01
    • 1970-01-01
    • 2021-09-30
    相关资源
    最近更新 更多