Cassandra迁移后UDT字段插入报CodecNotFoundException的解决建议
UDT变更导致CodecNotFoundException的解决方案
问题场景
在保持1.X版本应用Pod运行的前提下,将Cassandra(3.11)从1.X升级到1.Y版本,执行了UDT字段新增的ALTER脚本:
ALTER TYPE my_optional_params add refresh_context map <text,text>;
升级完成后,1.X版本应用向新数据库插入数据时抛出CodecNotFoundException,重启Pod可解决问题,但无法接受每次迁移都重启Pod的操作。
相关代码片段:
原UDT与表创建脚本
CREATE TYPE IF NOT exists my_optional_params ( validity_period map <text,text>, refresh ascii, my_flag boolean ); CREATE TABLE IF NOT exists table_X ( id timeuuid, name text, description text, owner text, optional_params frozen< my_optional_params >, PRIMARY KEY ((tenant_name), id) ) WITH gc_grace_seconds = 86400 AND comment = 'table X';
自定义TypeCodec核心代码
public class MyOptionalParamsTypeCodec implements TypeCodec<MyOptionalParams> { static final String VP = "validity_period"; static final String KEY_REFRESH = "refresh"; static final String KEY_MYFLAG = "my_flag"; private final TypeCodec<UdtValue> innerCodec; private final UserDefinedType userType; private final DataType cqlType; private final GenericType<MyOptionalParams> javaType; // 构造函数、转换逻辑等省略... }
异常信息
com.datastax.oss.driver.api.core.type.codec.CodecNotFoundException: Codec not found for requested operation: [UDT(my_keyspace.my_optional_params) <- -> com.api.model.MyOptionalParams]
原因分析
- 驱动元数据缓存不更新:Datastax驱动会缓存Cassandra的Schema元数据(包括UDT结构),UDT变更后,客户端缓存的旧元数据与新UDT结构不匹配,导致驱动无法找到对应Codec。
- 自定义Codec依赖旧元数据:自定义
MyOptionalParamsTypeCodec初始化时绑定了旧的UDT元数据,无法适配新增字段后的UDT类型。 - PreparedStatement缓存:DAO中缓存的
PreparedStatement基于旧UDT元数据生成,无法兼容新的UDT结构。
解决方案
1. 修改自定义Codec,适配UDT结构变化
更新Codec实现,使其不依赖初始化时的旧UDT元数据,每次转换都使用最新的UDT结构:
@Override public ByteBuffer encode( @Nullable MyOptionalParams myOptionalParams, @NonNull ProtocolVersion protocolVersion) { // 每次获取最新的UDT元数据 UserDefinedType latestUserType = connection.getCqlSession().getMetadata() .getKeyspace("my_keyspace") .flatMap(ks -> ks.getUserDefinedType("my_optional_params")) .orElseThrow(() -> new IllegalStateException("UDT my_optional_params not found")); UdtValue udtValue = latestUserType.newValue(); // 仅设置Java实体中存在的字段,忽略新增字段 if (myOptionalParams != null) { if (myOptionalParams.getValidityPeriod() != null) { Map<String, String> vpMap = new HashMap<>(); vpMap.put("unit", myOptionalParams.getValidityPeriod().getMode().name()); vpMap.put("value", myOptionalParams.getValidityPeriod().getValue().toString()); udtValue.setMap(VP, vpMap, String.class, String.class); } udtValue.setString(KEY_REFRESH, Optional.ofNullable(myOptionalParams.getRefresh()).orElse("")); // 检查字段是否存在,避免UDT变更后报错 if (latestUserType.contains(KEY_MYFLAG)) { udtValue.setBoolean(KEY_MYFLAG, myOptionalParams.isMyFlag()); } } return this.innerCodec.encode(udtValue, protocolVersion); }
同时修改解码逻辑,忽略新增的UDT字段:
Function<UdtValue, MyOptionalParams> udtToOptionalParamsTransformer = (udtValue) -> { MyOptionalParams myOptionalParams = new MyOptionalParams(); if (udtValue != null) { Map<String, String> vpInfo = udtValue.getMap(VP, String.class, String.class); if (vpInfo != null && !vpInfo.isEmpty()) { final String vpUnit = vpInfo.get("unit"); final Integer vpValue = Integer.valueOf(vpInfo.get("value")); myOptionalParams.setValidityPeriod(new ValidityPeriod(VPMode.valueOf(vpUnit), vpValue)); } myOptionalParams.setRefresh(udtValue.getString(KEY_REFRESH)); if (udtValue.getType().contains(KEY_MYFLAG)) { myOptionalParams.setMyFlag(udtValue.getBoolean(KEY_MYFLAG)); } // 忽略新增的refresh_context字段 } return myOptionalParams; };
2. 刷新驱动Schema元数据缓存
在UDT变更完成后,手动触发驱动的Schema刷新,获取最新的元数据:
// 在迁移脚本执行完成后调用 connection.getCqlSession().refreshSchema();
可以将此操作集成到数据库迁移流程中,或者提供一个内部接口用于手动触发刷新。
3. 避免PreparedStatement本地缓存
修改DAO中getInsertPreparedStatement方法,移除本地缓存,每次都重新生成PreparedStatement(若担心性能,可添加基于Schema版本的缓存刷新机制):
PreparedStatement getInsertPreparedStatement() { return this.prepare( QueryBuilder.insertInto("my_keyspace", "table_X") .value("id", QueryBuilder.bindMarker()) .value("name", QueryBuilder.bindMarker()) .value("description", QueryBuilder.bindMarker()) .value("owner", QueryBuilder.bindMarker()) .value("optional_params", QueryBuilder.bindMarker())); }
4. 先兼容再升级的迁移策略
采用应用先兼容,数据库后升级的流程:
- 先部署修改后的1.X版本应用,使其能同时兼容新旧UDT结构(即上述修改后的Codec)。
- 再执行UDT字段新增的ALTER脚本。
此方案可避免数据库升级期间的应用异常,无需重启Pod。
总结
通过修改自定义Codec适配UDT结构变化、刷新驱动元数据缓存、调整PreparedStatement缓存策略,或采用先兼容再升级的迁移流程,可彻底解决UDT变更后的CodecNotFoundException,无需重启应用Pod。
内容的提问来源于stack exchange,提问作者Vishant Vats
相关产品推荐
相关产品推荐

