Flink-Connector-Cassandra Schema解析错误:无法解析system_distributed.duration
问题背景
我们的Flink应用从Kafka消费数据并写入Scylla DB,启动时触发Schema解析错误。
环境版本
- Flink版本:1.15.4
- Scylla DB版本:scylla-enterprise-2022.2.11-0
依赖配置
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-cassandra_2.12</artifactId> <version>1.15.4</version> </dependency>
报错日志
WARNING: All illegal access operations will be denied in a future release 2023-08-11 09:16:04.567 [Source: Custom Source -> Map -> Sink: Cassandra Sink (7/9)#0] ERROR com.datastax.driver.core.SchemaParser - Error parsing schema for table system_distributed.service_levels: Cluster.getMetadata().getKeyspace("system_distributed").getTable("service_levels") will be missing or incomplete com.datastax.driver.core.exceptions.UnresolvedUserTypeException: Cannot resolve user type system_distributed.duration at com.datastax.driver.core.DataTypeCqlNameParser.parse(DataTypeCqlNameParser.java:147) at com.datastax.driver.core.TableMetadata.build(TableMetadata.java:188) at com.datastax.driver.core.SchemaParser.buildTables(SchemaParser.java:176) at com.datastax.driver.core.SchemaParser.buildKeyspaces(SchemaParser.java:128) at com.datastax.driver.core.SchemaParser.refresh(SchemaParser.java:64) at com.datastax.driver.core.ControlConnection.refreshSchema(ControlConnection.java:341) at com.datastax.driver.core.ControlConnection.tryConnect(ControlConnection.java:271) at com.datastax.driver.core.ControlConnection.reconnectInternal(ControlConnection.java:199) at com.datastax.driver.core.ControlConnection.connect(ControlConnection.java:77) at com.datastax.driver.core.Cluster$Manager.init(Cluster.java:1414)
原因分析
错误核心在于:Scylla系统库system_distributed下的service_levels表依赖自定义用户类型duration,但Flink Cassandra Connector依赖的DataStax Java Driver版本(Flink 1.15.4对应的Connector依赖Driver 3.x系列)在扫描系统库Schema时,无法正确识别或加载这个Scylla扩展的自定义类型,最终导致Schema解析失败。
解决方案建议
限制Schema扫描范围:
配置Cassandra Driver仅加载业务相关的Keyspace,避免扫描系统Keyspace。可以在Flink CassandraSink的配置中明确指定业务Keyspace,或者添加Driver配置参数datastax-java-driver.schema.enabled=false禁用自动Schema加载(注意:若业务表使用自定义类型,需手动处理类型映射)。升级依赖版本:
排查是否有兼容的Flink Cassandra Connector更新版本,或者手动替换为适配Scylla的DataStax Driver版本(需确保Driver版本与Flink Connector、Scylla版本兼容)。检查Scylla集群状态:
确认Scylla集群中system_distributed库的duration用户类型是否存在,若因集群同步问题导致缺失,可重启相关节点或手动同步Schema(系统表的类型通常由Scylla自动维护,优先排查节点健康状态)。
内容的提问来源于stack exchange,提问作者Jaiprasad

