【问题标题】:update one column in cassandra table更新 cassandra 表中的一列
【发布时间】:2016-08-09 05:28:24
【问题描述】:

我有一个 cassandra 表 person_master (personId: int, 客户 ID:整数, 名字:字符串, 姓氏:字符串, mrids: 设置)primaryKey(personId 和 customerID)

假设我有一个结构 [personId, customerId, firstName, lastname, messageType: String, source: String, sourceType: String] 的输入 RDD

假设RDD的值:[1001,119,None,None,{abc.xyz} 并且 cassandra 行的值为 [1001,119,Vikash,Singh,{aaa.bbb}]

我想根据 RDD 值获取 cassandra 行并更新 cassandra 表的 mrids 列并使用 cassandra 行中的所有其他列。

例如在此我希望最终的 RDD 值为 [1001,119,Vikash,Singh,{aaa.bbb,abc.xyz}],稍后我将更新为 cassandra。

谁能给我使用 cassandra 连接器在 Spark 中执行此操作的解决方案。

【问题讨论】:

    标签: scala apache-spark cassandra apache-spark-sql spark-cassandra-connector


    【解决方案1】:

    假设 sc 是 sparkContext 之类的,

    val sparkConf = new SparkConf().setMaster(SPARK_MASTER)
                                .setAppName(SPARK_SCALA_APP_NAME)
                                .setJars(SPARK_SCALA_JAR)
    sparkConf.set("spark.cassandra.connection.host", value)
    sparkConf.set("spark.cassandra.auth.username", value)
    sparkConf.set("spark.cassandra.auth.password", value)
    val sc = new SparkContext(sparkConf)
    

    可以使用或忽略where子句(where只能使用它的分区键)

    val selectedRow = sc.cassandraTable("keyspace", "tableName")
          .select("key", "column2", "column3")
          .where("key IN ?", keys)
          .as((key: String, column2: String, column3: Integer)
              =>(key, column2, column3))
    

    对你的 rdd 进行过滤和修改 然后像这样保存,

    selectedRow.saveToCassandra("keyspace",
                               "tableName",
                               SomeColumns("key", "column2", "column3"))
    

    【讨论】:

      猜你喜欢
      • 2013-02-19
      • 2013-04-27
      • 2017-04-22
      • 2015-10-22
      • 2016-10-14
      • 1970-01-01
      • 2019-08-30
      • 2019-12-10
      • 2015-10-04
      相关资源
      最近更新 更多