【问题标题】:How to update specific set of Cassandra columns from Spark Dataframe using Datastax connector如何使用 Datastax 连接器从 Spark Dataframe 更新特定的 Cassandra 列集
【发布时间】:2023-03-26 14:42:01
【问题描述】:

我有一个包含几列的 Cassandra 表,我想从 Spark 2.4.0 更新其中的一个(以及多列的内容?)。但是,如果我不提供所有列,则记录不会得到更新。

Cassandra 架构:

rowkey,message,number,timestamp,name
1,hello,12345,12233454,ABC

重点是 Spark DataFrame 由 rowkey 和必须在 Cassandra 表中更新的更新时间戳组成。

我尝试在选项之后选择列,但似乎没有这样的方法。

finalDF.select("rowkey","current_ts")
  .withColumnRenamed("current_ts","timestamp")
  .write
  .format("org.apache.spark.sql.cassandra")
  .options(Map("table" -> "table_data", "keyspace" -> "ks_data"))
  .mode("overwrite")
  .option("confirm.truncate","true")
  .save()

说,

finalDF=
rowkey,current_ts
1,12233999

那么 Cassandra 表应该保持更新后的值,

rowkey,message,number,timestamp,name
1,hello,12345,12233999,ABC

我正在使用 Dataframe API。所以不能使用rdd方法。我怎么能做到这一点? Cassandra 版本 3.11.3,Datastax 连接器 2.4.0-2.11

【问题讨论】:

  • 因此将 Savemode 更改为“append”解决了这个问题。有什么说明吗?

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


【解决方案1】:

澄清SaveMode 用于指定将DataFrame 保存到数据源的预期行为。(不仅适用于c*,也适用于任何数据源)。可用options 有

  1. SaveMode.ErrorIfExists
  2. SaveMode.Append
  3. SaveMode.Overwrite
  4. SaveMode.Ignore

在这种情况下,由于您已经有数据并且想要追加,您必须使用SaveMode.Append

import org.apache.spark.sql.SaveMode

finalDF.select("rowkey","current_ts")
  .withColumnRenamed("current_ts","timestamp")
  .write
  .format("org.apache.spark.sql.cassandra")
  .options(Map("table" -> "table_data", "keyspace" -> "ks_data"))
  .mode(SaveMode.Append)
  .option("confirm.truncate","true")
  .save()

在 SaveModes 上查看 spark 文档

【讨论】:

    猜你喜欢
    • 2018-03-01
    • 2015-05-24
    • 2017-02-13
    • 2016-08-11
    • 2015-10-28
    • 2015-08-16
    • 2017-03-04
    • 2020-10-17
    • 2016-02-04
    相关资源
    最近更新 更多