Flink CDC YAML配置SQL Server间同步报错:找不到sqlserver-cdc工厂
问题排查(YAML同步方式)
1. 依赖包完整性检查
你已将flink-sql-connector-sqlserver-cdc-3.1.0.jar放入/flink-1.19.1/lib/,但需补充以下验证:
- 解压
flink-cdc-3.1.0-bin.tar.gz后,需将包内flink-cdc-base相关jar也复制到lib目录,flink-cdc.sh脚本依赖这些基础组件加载CDC工厂类。 - 执行
jar tf flink-sql-connector-sqlserver-cdc-3.1.0.jar查看包内是否存在SqlServerCDCSourceFactory类,该类是sqlserver-cdc标识符对应的核心工厂类,缺失则会触发报错。
2. YAML配置修正(参考MySQL格式适配)
官方无SQL Server的YAML示例,可基于JDBC Sink+SQL Server CDC Source编写配置,核心内容如下:
source: type: sqlserver-cdc hostname: 源SQL Server地址 port: 1433 username: 数据库用户名 password: 数据库密码 database-name: TEST_FOR_FLINK table-name: dbo.orders server-time-zone: Asia/Shanghai sink: type: jdbc url: jdbc:sqlserver://目标SQL Server地址:1433;databaseName=DW username: 数据库用户名 password: 数据库密码 table-name: dbo.orders_dw driver-class-name: com.microsoft.sqlserver.jdbc.SQLServerDriver pipeline: name: SQLServer-To-SQLServer-Sync
3. 脚本执行路径校验
执行命令时需确保当前目录为/flink-1.19.1,或使用绝对路径执行:
cd /flink-1.19.1 ./flink-cdc.sh pipeline/mssql-to-mssql-test01.yaml
同时检查flink-cdc.sh脚本中的classpath配置,确认是否包含lib目录下所有jar包,避免脚本未加载到依赖。
Java + Table API 同步Demo
如果YAML方式暂时无法解决,可采用Table API实现稳定同步,步骤如下:
1. Maven依赖配置(适配你的版本)
在项目pom.xml中添加以下依赖:
<dependencies> <!-- Flink Table API基础依赖 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-api-java</artifactId> <version>1.19.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-planner-loader</artifactId> <version>1.19.1</version> <scope>provided</scope> </dependency> <!-- SQL Server CDC Source --> <dependency> <groupId>com.ververica</groupId> <artifactId>flink-sql-connector-sqlserver-cdc</artifactId> <version>3.1.0</version> </dependency> <!-- SQL Server JDBC Sink --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-jdbc</artifactId> <version>3.1.2-1.18</version> </dependency> <dependency> <groupId>com.microsoft.sqlserver</groupId> <artifactId>mssql-jdbc</artifactId> <version>12.6.3.jre11</version> </dependency> </dependencies>
2. 同步代码实现
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; public class SqlServerSyncJob { public static void main(String[] args) throws Exception { // 初始化执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env); // 创建源表(SQL Server CDC) tableEnv.executeSql("CREATE TABLE orders_source (" + " order_id INT," + " customer_id INT," + " order_date TIMESTAMP(3)," + " amount DECIMAL(10,2)," + " PRIMARY KEY(order_id) NOT ENFORCED" + ") WITH (" + " 'connector' = 'sqlserver-cdc'," + " 'hostname' = '源SQL Server地址'," + " 'port' = '1433'," + " 'username' = '用户名'," + " 'password' = '密码'," + " 'database-name' = 'TEST_FOR_FLINK'," + " 'table-name' = 'dbo.orders'," + " 'server-time-zone' = 'Asia/Shanghai'" + ")"); // 创建目标表(JDBC Sink) tableEnv.executeSql("CREATE TABLE orders_sink (" + " order_id INT," + " customer_id INT," + " order_date TIMESTAMP(3)," + " amount DECIMAL(10,2)," + " PRIMARY KEY(order_id) NOT ENFORCED" + ") WITH (" + " 'connector' = 'jdbc'," + " 'url' = 'jdbc:sqlserver://目标SQL Server地址:1433;databaseName=DW'," + " 'username' = '用户名'," + " 'password' = '密码'," + " 'table-name' = 'dbo.orders_dw'," + " 'driver' = 'com.microsoft.sqlserver.jdbc.SQLServerDriver'" + ")"); // 执行全量+增量同步 tableEnv.executeSql("INSERT INTO orders_sink SELECT * FROM orders_source"); env.execute("SQLServer-Sync-Job"); } }
3. 打包与运行
- 用Maven将项目打包为fat jar(可通过
maven-shade-plugin处理依赖冲突)。 - 将jar上传到CentOS的
/flink-1.19.1目录,执行启动命令:
./bin/flink run -c SqlServerSyncJob your-jar-file-name.jar
内容的提问来源于stack exchange,提问作者潘德拉贡阿尔托莉雅
相关产品推荐
相关产品推荐

