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
相关产品推荐
相关产品推荐

