基于Spark和Scala将扁平文件转换为Cassandra数据模型
我之前也碰到过类似的嵌套Frozen类型写入Cassandra的棘手问题,完全不用RDD就能搞定——核心是靠Scala case class映射Cassandra的Frozen UDT,再通过DataFrame的聚合转换生成符合结构的数据,最后用Spark Cassandra Connector完成写入。下面是一步步的实操方案:
1. 定义对应Cassandra UDT的Scala Case Class
首先得给Cassandra里的forder和flineitem这两个Frozen UDT,定义匹配的Scala case class,注意字段名要和Cassandra的列名完全对齐(Cassandra默认小写,所以case class字段也用小写更稳妥):
// 对应Cassandra的forder Frozen UDT case class forder(orderNumber: String) // 对应Cassandra的flineitem Frozen UDT case class flineitem(lineitemid: Int, trackingnumber: String) // 对应Cassandra的Transfer表整体结构 case class Transfer(transferNumber: String, orderList: List[forder], lineItem: List[flineitem])
提示:如果你的UDT里还有更多字段,直接在对应的case class里追加就行,类型要和Cassandra的类型严格对应(比如TEXT对应String,INT对应Int)。
2. 把扁平文件DataFrame转成嵌套结构
假设你读入扁平文件后得到的DataFrame是这样的(每行对应一条原始明细):
| transferNumber | orderNumber | lineitemid | trackingnumber |
|---|---|---|---|
| T001 | O001 | 1 | TRK001 |
| T001 | O002 | 2 | TRK002 |
| T002 | O003 | 3 | TRK003 |
你需要通过分组+聚合,把同一个transferNumber下的明细数据收拢成嵌套列表。用groupBy搭配agg,结合collect_list和struct就能实现:
import org.apache.spark.sql.functions._ // flatDF是读入扁平文件后的原始DataFrame val nestedDF = flatDF .groupBy("transferNumber") .agg( // 把orderNumber打包成forder结构体,再收集成List collect_list(struct("orderNumber").as[forder]).alias("orderList"), // 把lineitemid和trackingnumber打包成flineitem结构体,再收集成List collect_list(struct("lineitemid", "trackingnumber").as[flineitem]).alias("lineItem") )
这里的
struct负责把零散字段组合成对应UDT的结构体,.as[CaseClass]强制转换成我们定义的case class类型,collect_list则把同分组下的结构体收拢成List,完美匹配Cassandra的Frozen List类型。
3. 写入Cassandra
先确保你已经引入了对应版本的Spark Cassandra Connector依赖(比如Spark 3.3.x对应com.datastax.spark:spark-cassandra-connector_2.12:3.3.0),然后配置连接参数写入数据:
nestedDF.write .format("org.apache.spark.sql.cassandra") .options(Map( "keyspace" -> "你的keyspace名称", "table" -> "transfer", "spark.cassandra.connection.host" -> "Cassandra节点地址", "spark.cassandra.connection.port" -> "9042" )) .mode(SaveMode.Append) // 根据需求选Append/Overwrite等模式 .save()
常见问题排查
- 类型不匹配报错:检查case class字段类型和Cassandra列类型是否严格对应,比如Cassandra的INT不能用Scala的Long代替,TEXT对应String别写错。
- 字段名不匹配:Cassandra默认列名小写,case class字段尽量保持一致;如果要忽略大小写,可以加配置
"spark.cassandra.mapping.caseSensitive" -> "false"。 - 版本兼容问题:Spark版本和Cassandra Connector版本必须匹配,比如Spark 3.2.x对应Connector 3.2.x,版本不对会出现各种奇怪的序列化错误。
- 序列化失败:嵌套的case class要定义成顶级类(别放在函数内部),case class默认是可序列化的,但内部类可能会出问题。
内容的提问来源于stack exchange,提问作者Vigneshwaran

