Spark 3 Cassandra Connector读取含空UDT列数据时遇类型转换错误
解决Spark 3.x读取Cassandra含Null值UDT列的类型转换异常
问题背景
从Cassandra表读取数据时,遇到phone_numbers(UDT类型列)存在Null值的情况。在Spark 2.4.8搭配Spark Cassandra Connector 2.5.1环境下代码可正常运行,但升级至Spark 3.2.2 + Connector 3.2.0后,抛出如下异常:
com.datastax.spark.connector.types.TypeConversionException: Cannot convert object null to example.PhoneNumbers
数据表结构:
| account_id | phone_numbers |
|---|---|
| 1 | {mobile_phone: {phone_number: 1234567890, date_added: 2022-01-01}, home_phone: {phone_number: 1234567891, date_added: 2022-01-02}} |
| 2 | null |
原读取代码:
CassandraJavaRDD<PhoneNumbers> cassandraJavaRDD = CassandraJavaUtil.javaFunctions(javaSparkContext) .cassandraTable( keyspace, tableName, mapRowTo(PhoneNumbers.class)) .select(columns); Dataset<PhoneNumbers> tableAsDataset = sparkSession.createDataset(cassandraJavaRDD.rdd(), Encoders.bean(PhoneNumbers.class)); tableAsDataset.show();
解决方案
1. 用Optional包装UDT类型
修改实体类,将PhoneNumbers字段替换为Optional<PhoneNumbers>,让连接器自动处理Null值映射:
// 示例Account实体类 public class Account { private Integer accountId; private Optional<PhoneNumbers> phoneNumbers; // 省略getter、setter方法 }
读取时映射到Account类:
CassandraJavaRDD<Account> cassandraJavaRDD = CassandraJavaUtil.javaFunctions(javaSparkContext) .cassandraTable( keyspace, tableName, mapRowTo(Account.class)) .select("account_id", "phone_numbers"); Dataset<Account> tableAsDataset = sparkSession.createDataset(cassandraJavaRDD.rdd(), Encoders.bean(Account.class)); tableAsDataset.show();
2. 自定义类型转换器处理Null值
实现TypeConverter接口,显式处理Null值(比如返回空对象):
public class PhoneNumbersConverter extends TypeConverter<PhoneNumbers> { @Override public PhoneNumbers convert(Object obj) { if (obj == null) { // 根据业务需求返回空的PhoneNumbers实例 return new PhoneNumbers(); } // 复用连接器的UDT转换逻辑处理非Null值 UDTValue udtValue = (UDTValue) obj; PhoneNumbers phoneNumbers = new PhoneNumbers(); phoneNumbers.setMobilePhone(udtValue.getUDTValue("mobile_phone", mapRowTo(Phone.class))); phoneNumbers.setHomePhone(udtValue.getUDTValue("home_phone", mapRowTo(Phone.class))); return phoneNumbers; } }
注册转换器后再读取数据:
// 注册自定义转换器 TypeConverter.registerConverter(PhoneNumbers.class, new PhoneNumbersConverter()); // 原有读取代码保持不变 CassandraJavaRDD<PhoneNumbers> cassandraJavaRDD = CassandraJavaUtil.javaFunctions(javaSparkContext) .cassandraTable( keyspace, tableName, mapRowTo(PhoneNumbers.class)) .select(columns);
3. 提前过滤含Null值的行
如果业务允许忽略Null行,可以在读取阶段直接过滤:
- RDD API过滤:
CassandraJavaRDD<PhoneNumbers> cassandraJavaRDD = CassandraJavaUtil.javaFunctions(javaSparkContext) .cassandraTable( keyspace, tableName, mapRowTo(PhoneNumbers.class)) .select(columns) .filter(Objects::nonNull);
- CQL条件过滤(更高效,在Cassandra端过滤):
CassandraJavaRDD<PhoneNumbers> cassandraJavaRDD = CassandraJavaUtil.javaFunctions(javaSparkContext) .cassandraTable( keyspace, tableName, mapRowTo(PhoneNumbers.class)) .where("phone_numbers IS NOT NULL");
4. 配置连接器空值处理参数
针对DataFrame API场景,可配置spark.cassandra.sql.nullToNone参数(默认true),辅助处理Null值映射:
sparkSession.conf().set("spark.cassandra.sql.nullToNone", "true");
如果使用RDD API,建议结合前面的方法一起使用。
内容的提问来源于stack exchange,提问作者Austin Spencer
相关产品推荐
相关产品推荐

