Flink 1.7 PyFlink对接MongoDB转MySQL遇ClassCastException求助
你的ClassCastException本质是版本不兼容导致的类加载冲突,具体是Flink版本、MongoDB连接器版本、MongoDB Java驱动版本三者完全不匹配,以下是具体解决步骤:
1. 替换适配Flink 1.7的MongoDB连接器
你当前用的flink-sql-connector-mongodb-1.1.0-1.17.jar是给Flink 1.17版本设计的,和Flink 1.7完全不兼容。Flink 1.7没有官方SQL MongoDB连接器,需要使用对应版本的第三方连接器,比如flink-connector-mongodb_2.11-1.7.2.jar(注意匹配你的Scala版本,Flink 1.7默认使用Scala 2.11)。
2. 匹配MongoDB驱动版本
不要使用5.1.2这类高版本MongoDB驱动,Flink 1.7适配的是MongoDB 3.x系列驱动,比如mongodb-driver-3.8.2.jar、mongodb-driver-core-3.8.2.jar、bson-3.8.2.jar。过高的驱动版本会和老版本Flink的类加载机制冲突,导致同一个BsonDocument类被不同类加载器加载,JVM会判定为不同类型,从而触发类型转换异常。
3. 清理冲突Jar包
- 删除你当前添加的
flink-sql-connector-mongodb-1.1.0-1.17.jar、bson-5.1.2.jar、mongodb-driver-sync-5.1.2.jar、mongodb-driver-core-5.1.2.jar - 将适配的Jar包放入Flink的
lib目录,或者通过env.add_jars指定正确的文件路径
4. 调整表创建语句
适配Flink 1.7的MongoDB连接器参数和新版本有所不同,示例如下:
t_env.execute_sql(""" CREATE TABLE source ( id STRING, name STRING, -- 其他字段定义 ) WITH ( 'connector.type' = 'mongodb', 'connector.hosts' = 'mongodb://localhost:27017', 'connector.database' = '你的数据库名', 'connector.collection' = '你的集合名', 'connector.username' = '用户名', 'connector.password' = '密码' ) """)
问题根源说明
报错中提示org.bson.BsonDocument无法赋值给同类型字段,这是类加载器隔离导致的:高版本驱动与连接器中的BsonDocument类被不同类加载器加载,即使全类名相同,JVM也会判定为不同类型。Flink 1.7的类加载机制与新版本差异较大,高版本组件无法兼容。
内容的提问来源于stack exchange,提问作者Chun Lih Chiang

