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

Scala Spark使用UDF转换Cassandra 4000万数据写入卡顿,求优化方案

Scala Spark 写入Cassandra卡顿优化方案

当前4000万数据迁移卡顿主要由以下几个原因导致:

  • 自定义UDF处理JSON效率极低,每条数据重复初始化序列化工具,无向量化优化
  • 多余的count、persist操作导致全表扫描次数翻倍,内存占用过高
  • Cassandra连接器参数配置不合理,写入批次、并发数不适配场景
  • 本地运行模式资源上限低,不适合处理千万级数据量

1. 替换自定义UDF为Spark内置JSON函数(性能提升5~10倍)

Spark内置的JSON解析、构造函数是向量化执行的,远快于逐行处理的自定义UDF,完全可以实现你的转换逻辑:

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

// 定义原data字段的schema
val dataSchema = StructType(Seq(
  StructField("type", IntegerType),
  StructField("actionData", StructType(Seq(
    StructField("actionType", StringType),
    StructField("id", StringType),
    StructField("title", StringType),
    StructField("description", StringType),
    StructField("config", StructType(Seq(StructField("type", StringType)))
  )))
))

// 构造action字段,完全不用UDF
val userFeedV2Format = userFeedData
  .withColumn("id", expr("uuid()"))
  .withColumn("data_json", from_json(col("data"), dataSchema))
  // 提取需要的字段构造template
  .withColumn("template", struct(
    to_json(struct(col("data_json.actionData.title"), col("data_json.actionData.description"))).alias("data"),
    lit("4.4.0").alias("ver"),
    lit("JSON").alias("type")
  ))
  // 构造additionalInfo
  .withColumn("additionalInfo", struct(
    col("data_json.type").alias("type"),
    col("data_json.actionData.id").alias("id"),
    col("data_json.actionData.config").alias("config")
  ))
  // 构造最终的action字段
  .withColumn("action", to_json(struct(
    col("data_json.actionData.actionType").alias("type"),
    col("category"),
    struct(lit("System").alias("type"), col("createdby").alias("id")).alias("createdBy"),
    col("template"),
    col("additionalInfo")
  )))
  .select(
    col("id"),col("category"),col("createdby"),
    col("createdon"),col("action"),col("expireon"),
    col("priority"),col("status"),col("updatedby"),
    col("updatedon"),col("userid"),
    lit("v2").cast(StringType).alias("version")
  )

2. 去掉冗余操作,减少作业触发次数

  • 删掉所有中间persist()和unpersist()代码,你的流程中所有数据集只被使用一次,持久化完全是多余的内存开销
  • 删掉中间两次count()操作,每次count都会触发一次全表扫描,4000万条数据等于多跑了两次完整流程,需要校验的话可以写完后再单独统计数量

3. 调整Spark和Cassandra连接器参数

启动命令调整

如果必须用本地模式运行,修改启动参数:

bin/spark-shell --master local[*] --packages com.datastax.spark:spark-cassandra-connector_2.11:2.5.0 --driver-memory 200g --conf spark.sql.shuffle.partitions=64

本地模式下shuffle分区数设为CPU核心数的2~3倍即可,默认200会导致小任务过多。

代码中配置调整

SparkSession
  .builder()
  .appName("UserLookupMigration")
  .config("spark.master", "local[*]")
  .config("spark.cassandra.connection.host",cassandra)
  // 写入批次调整为1000,避免过大导致Cassandra节点超时
  .config("spark.cassandra.output.batch.size.rows", "1000")
  .config("spark.cassandra.read.timeoutMS", "120000")
  .config("spark.cassandra.connection.timeoutMS", "60000")
  // 写入并发调整为200,根据Cassandra集群负载可上下调整
  .config("spark.cassandra.output.concurrent.writes", "200")
  .config("spark.cassandra.output.consistency.level", "LOCAL_QUORUM")
  .getOrCreate()

4. 写入前调整分区数

写入Cassandra前调整分区数,避免分区过多导致连接数超限,也避免分区过少导致并行度不足:

// 4000万条数据建议分200~400个分区写入,可根据实际资源调整
userFeedV2Format.repartition(200)
  .write.format("org.apache.spark.sql.cassandra")
  .option("keyspace", "sunbird_notifications")
  .option("table", "TABLE2")
  .mode(SaveMode.Append).save()

如果必须保留自定义UDF,将ObjectMapper提取为单例,避免每条数据重复初始化:

object JsonUtil extends Serializable {
  val mapper: ObjectMapper = new ObjectMapper().registerModule(DefaultScalaModule)
}

// UDF中直接用JsonUtil.mapper,不要每次new
def myColUpdate= udf((data: String, createdby: String, category: String)=> {
  // 原有逻辑,把new ObjectMapper换成JsonUtil.mapper
})

优化完成后千万级数据的写入效率会提升10倍以上,卡顿问题基本可以解决。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 04:45:04