You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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]

原因分析

  1. 驱动元数据缓存不更新:Datastax驱动会缓存Cassandra的Schema元数据(包括UDT结构),UDT变更后,客户端缓存的旧元数据与新UDT结构不匹配,导致驱动无法找到对应Codec。
  2. 自定义Codec依赖旧元数据:自定义MyOptionalParamsTypeCodec初始化时绑定了旧的UDT元数据,无法适配新增字段后的UDT类型。
  3. 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. 先部署修改后的1.X版本应用,使其能同时兼容新旧UDT结构(即上述修改后的Codec)。
  2. 再执行UDT字段新增的ALTER脚本。
    此方案可避免数据库升级期间的应用异常,无需重启Pod。

总结

通过修改自定义Codec适配UDT结构变化、刷新驱动元数据缓存、调整PreparedStatement缓存策略,或采用先兼容再升级的迁移流程,可彻底解决UDT变更后的CodecNotFoundException,无需重启应用Pod。

内容的提问来源于stack exchange,提问作者Vishant Vats

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 13:37:04