Spark写入私有属性类RDD到Cassandra列名不匹配解决方法
实现方案
不需要修改原有类的私有属性、属性名,也无需将类改为Case Class,通过自定义Spark Cassandra Connector的列映射规则即可完成写入。
核心原理
默认的DefaultColumnMapper会直接通过反射读取类的公有字段做匹配,既无法识别私有属性,也不会自动完成驼峰命名到蛇形命名的转换,因此会报字段找不到的异常。当前实体类是标准JavaBean风格(带完整getter/setter),直接使用JavaBeanColumnMapper即可通过getter方法读取属性值,同时支持自定义字段映射规则,完全满足约束要求。
具体实现
- 导入映射相关依赖包
import com.datastax.spark.connector.mapper.{ColumnMapper, JavaBeanColumnMapper} import org.apache.spark.SparkContext import com.datastax.spark.connector.cql.CassandraConnector import org.apache.spark.rdd.RDD
- 在写入逻辑的作用域内,为
Test类定义隐式列映射,手动绑定类属性(和getter命名对应)和Cassandra表字段的对应关系:
// 声明Test类和Cassandra表的字段映射关系 implicit val testColumnMapper: ColumnMapper[Test] = new JavaBeanColumnMapper[Test]( Map( "id" -> "id", "randomNumber" -> "random_number", "lastUpdater" -> "last_update" ) )
- 保留原有写入逻辑即可正常执行,不需要修改实体类代码:
implicit val connector = CassandraConnector(sc.getConf) val rdd: RDD[Test] = // 原有RDD生成逻辑 rdd.saveToCassandra("test", "test")
扩展优化
如果项目中有大量驼峰命名的JavaBean实体需要写入蛇形命名的Cassandra表,可以自定义一个通用的驼峰转蛇形映射器,避免每个类手动编写字段映射:
import java.util.Locale class CamelToSnakeJavaBeanMapper[T] extends JavaBeanColumnMapper[T]( columnNameOverride = Map.empty, // 自动将驼峰格式属性名转换为蛇形格式字段名 propertyToColumnName = (propName: String) => propName.replaceAll("([a-z])([A-Z])", "$1_$2").toLowerCase(Locale.ROOT) ) // 对应Test类只需要声明隐式Mapper即可,特殊不匹配字段可在columnNameOverride中单独配置 implicit val testMapper: ColumnMapper[Test] = new CamelToSnakeJavaBeanMapper[Test]()
注意:
columnNameOverride中配置的单字段映射优先级高于自动转换规则,遇到命名特殊的字段单独配置即可。
内容的提问来源于stack exchange,提问作者Alejandro Espada
相关产品推荐
相关产品推荐

