You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark DataFrame含重复主键时,如何向Cassandra写入最新行?

Cassandra与Spark Connector数据写入问题

场景说明

我们有一张Cassandra表person,建表语句如下:

CREATE TABLE test.person (
    name text PRIMARY KEY,
    score bigint
)

对应的Spark DataFrame数据如下:

val caseClassDF = Seq(Person("Andy1", 32), Person("Mark1", 27), Person("Ron", 27),Person("Andy1", 20),Person("Ron", 270),Person("Ron", 2700),Person("Mark1", 37),Person("Andy1", 200),Person("Andy1", 2000)).toDF()

DataFrame中存在多条主键(name)相同的记录,写入Cassandra时遇到以下疑问及现象:

测试现象

使用spark-cassandra-connector_2.12:3.2.1,写入代码如下:

val spark = SparkSession.builder()
    .master("local[1]")
    .appName("CassandraConnector")
    .config("spark.cassandra.connection.host", "")
    .config("spark.cassandra.connection.port", "")
    .config("spark.sql.extensions", "com.datastax.spark.connector.CassandraSparkExtensions")
    .getOrCreate()

val caseClassDF = Seq(Person("Andy1", 32), Person("Mark1", 27), Person("Ron", 27),Person("Andy1", 20),Person("Ron", 270),Person("Ron", 2700),Person("Mark1", 37),Person("Andy1", 200),Person("Andy1", 2000)).toDF()

caseClassDF.write
      .format("org.apache.spark.sql.cassandra")
      .option("keyspace", "test")
      .option("table", "person")
      .mode("APPEND")
      .save()
  1. 当配置.master("local[1]")时,Cassandra表中Andy1的score始终为2000,Ron的score始终为2700,即序列中的最后一条记录。
  2. 当配置改为.master("local[*]")或.master("local[2]")时,Cassandra表中Andy1的score为随机值(如200或32)。

注:每次测试均使用全新的表,所有插入和更新在一个批次完成。


问题解答

问题1:Cassandra Connector内部如何处理行的顺序?

Spark是分布式计算框架,当使用多核心(如local[2])时,DataFrame的分区会被并行处理,Connector无法保证写入Cassandra的行顺序与DataFrame原始顺序一致:

  • 单核心(local[1])模式下,数据按顺序串行处理,同主键的最后一条记录会覆盖之前的,因此保留序列末尾的值。
  • 多核心模式下,不同分区的同主键记录可能被并行写入Cassandra,最终保留哪条取决于Cassandra接收请求的顺序,而非DataFrame原始顺序,因此出现随机结果。

问题2:从Kafka读取数据写入Cassandra,批量数据中有多事件场景,如何确保保存最新的score?

核心思路是在写入Cassandra前对同主键数据去重,只保留最新记录,具体方案如下:

方案1:DataFrame分组聚合保留最新记录

如果每条Kafka事件带有时间戳(或可通过事件顺序推断版本),按主键分组后取时间戳最大的记录:

import org.apache.spark.sql.functions._

// 假设DataFrame自带event_time字段,若没有可从Kafka元数据获取或生成当前时间
val latestDF = caseClassDF
  .groupBy("name")
  .agg(max("event_time").alias("latest_time"), first("score").alias("latest_score"))
  .select("name", "latest_score")
  .withColumnRenamed("latest_score", "score")

latestDF.write
  .format("org.apache.spark.sql.cassandra")
  .option("keyspace", "test")
  .option("table", "person")
  .mode("APPEND")
  .save()

若无时间戳,但DataFrame顺序对应事件先后(如Kafka消费顺序),可添加递增序号后分组取最大序号的记录:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

// 按主键分区,倒序排列后取第一条(即原始序列最后一条)
val windowSpec = Window.partitionBy("name").orderBy(lit(1).desc)
val latestDF = caseClassDF
  .withColumn("row_num", row_number().over(windowSpec))
  .filter(col("row_num") === 1)
  .drop("row_num")

latestDF.write
  .format("org.apache.spark.sql.cassandra")
  .option("keyspace", "test")
  .option("table", "person")
  .mode("APPEND")
  .save()

方案2:使用Cassandra轻量级事务(LWT)

若无法在Spark层去重,可利用Cassandra的条件更新,前提是每条记录有递增版本号:

  1. 修改Cassandra表增加版本字段:
ALTER TABLE test.person ADD version bigint;
  1. 写入时仅当新版本号更大时更新:
import com.datastax.spark.connector._

caseClassDF.rdd.map(row => {
  val name = row.getAs[String]("name")
  val score = row.getAs[Long]("score")
  val version = row.getAs[Long]("version") // DataFrame需包含递增的version字段
  (name, score, version)
}).saveToCassandra("test", "person", SomeColumns("name", "score", "version"),
  writeConf = WriteConf(ifNotExists = false)
    .withCASCondition(Condition("version", LT, placeholder))
)

注:LWT会增加Cassandra性能开销,适合数据量不大或强一致性要求场景。

方案3:调整Spark写入配置(不推荐)

设置spark.cassandra.output.batch.size.rows为1,强制单条写入,但会严重降低性能,仅适合测试场景。


内容的提问来源于stack exchange,提问作者Mayuri

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 23:03:08