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()
- 当配置
.master("local[1]")时,Cassandra表中Andy1的score始终为2000,Ron的score始终为2700,即序列中的最后一条记录。 - 当配置改为
.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的条件更新,前提是每条记录有递增版本号:
- 修改Cassandra表增加版本字段:
ALTER TABLE test.person ADD version bigint;
- 写入时仅当新版本号更大时更新:
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
相关产品推荐
相关产品推荐

