已部署MySQL CDC依赖,PyFlink作业仍报找不到mysql-cdc工厂问题咨询
错误原因
Jar包冗余冲突
你同时将flink-cdc-dist-3.3.0.jar、flink-cdc-pipeline-connector-mysql-3.3.0.jar和flink-sql-connector-mysql-cdc-3.3.0.jar放入集群/opt/flink/lib目录,其中flink-cdc-dist是CDC全量依赖包,flink-sql-connector-mysql-cdc是专门的MySQL CDC SQL连接器包,二者功能重叠,会导致类加载冲突,使Flink无法正确识别mysql-cdc的工厂类。不必要的Jar包引入
flink-cdc-pipeline-connector-mysql是用于Pipeline模式的连接器,而你采用的是PyFlink SQL方式,这个包完全不需要,进一步加剧了类路径混乱。
正确处理步骤
1. 清理冗余Jar包
进入Docker容器的/opt/flink/lib目录,删除以下冗余Jar:
rm flink-cdc-dist-3.3.0.jar flink-cdc-pipeline-connector-mysql-3.3.0.jar
仅保留flink-sql-connector-mysql-cdc-3.3.0.jar和mysql-connector-java-8.0.27.jar即可,前者已包含MySQL CDC SQL连接器的所有核心依赖。
2. 正确的作业提交方式
由于Jar包已放入集群全局类路径/opt/flink/lib,直接用以下命令提交作业:
flink run -py flink_job.py
如果requirements.txt包含PyFlink依赖或自定义Python包,确保Docker容器内已安装这些依赖。若需传递Python依赖包,可保留-pyFiles参数:
flink run -py flink_job.py -pyFiles requirements.txt
若不想修改集群镜像,也可在提交时通过-j参数指定Jar包路径(适合临时测试):
flink run -py flink_job.py -j /path/to/flink-sql-connector-mysql-cdc-3.3.0.jar -j /path/to/mysql-connector-java-8.0.27.jar
3. 验证代码表定义
确保PyFlink代码中MySQL-CDC源表的连接器配置正确,示例如下:
from pyflink.table import TableEnvironment, EnvironmentSettings env_settings = EnvironmentSettings.in_streaming_mode() t_env = TableEnvironment.create(env_settings) # MySQL-CDC源表定义 t_env.execute_sql(""" CREATE TABLE mysql_source ( id INT, name STRING, PRIMARY KEY(id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 'mysql-host', 'port' = '3306', 'username' = 'root', 'password' = 'password', 'database-name' = 'test_db', 'table-name' = 'test_table' ) """) # S3 Sink表定义 t_env.execute_sql(""" CREATE TABLE s3_sink ( id INT, name STRING ) WITH ( 'connector' = 'filesystem', 'path' = 's3://your-bucket/path', 'format' = 'parquet', 'sink.partition-commit.delay' = '1 min', 'sink.partition-commit.policy.kind' = 'success-file' ) """) # 数据插入逻辑 t_env.execute_sql("INSERT INTO s3_sink SELECT id, name FROM mysql_source").wait()
内容的提问来源于stack exchange,提问作者Rinze

