Flink CDC 实战教程:SqlServer CDC 到 Elasticsearch 实时数据同步
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
本文是一篇基于 Apache Flink CDC 开源项目的完整实战指南,核心演示如何通过 Flink SQL 将 SQL Server 数据库中的变更数据(CDC)实时捕获并同步到 Elasticsearch,并结合 Kibana 可视化验证同步结果。读者将掌握 SQL Server CDC 的启用方法、sqlserver-cdc连接器的完整配置、基于主键的维表关联(LEFT JOIN)实时宽表构建,以及增量变更(INSERT/UPDATE/DELETE)在 Elasticsearch 中的实时反映,可直接复制本文环境与 SQL 用于本地实验。
1. 实验环境与整体架构
本教程通过 Docker Compose 一键拉起三套服务,构成一个完整的实时数据同步演示环境:
- SqlServer(
mcr.microsoft.com/mssql/server:2019-latest):作为 CDC 数据源,存放业务表orders与products; - Elasticsearch(
elastic/elasticsearch:7.6.0):作为同步目标,存储orders与products关联(JOIN)后的结果; - Kibana(
elastic/kibana:7.6.0):用于可视化查看 Elasticsearch 中的实时数据变化。
创建docker-compose.yml文件,内容如下:
version: '2.1' services: sqlserver: image: mcr.microsoft.com/mssql/server:2019-latest container_name: sqlserver ports: - "1433:1433" environment: - "MSSQL_AGENT_ENABLED=true" - "MSSQL_PID=Standard" - "ACCEPT_EULA=Y" - "SA_PASSWORD=Password!" elasticsearch: image: elastic/elasticsearch:7.6.0 container_name: elasticsearch environment: - cluster.name=docker-cluster - bootstrap.memory_lock=true - "ES_JAVA_OPTS=-Xms512m -Xmx512m" - discovery.type=single-node ports: - "9200:9200" - "9300:9300" ulimits: memlock: soft: -1 hard: -1 nofile: soft: 65536 hard: 65536 kibana: image: elastic/kibana:7.6.0 container_name: kibana ports: - "5601:5601" volumes: - /var/run/docker.sock:/var/run/docker.sock几个关键配置说明:
MSSQL_AGENT_ENABLED=true必须开启,SQL Server CDC 依赖 SQL Server Agent 进程扫描事务日志并写入变更表;SA_PASSWORD=Password!为sa超级管理员密码,后续 Flink SQL 中username/password与此保持一致;- Elasticsearch 以单节点(
discovery.type=single-node)模式运行,ES_JAVA_OPTS限制堆内存为 512MB,避免实验机器资源吃紧。
在包含docker-compose.yml的目录下执行以下命令启动全部容器(后台守护模式):
docker-compose up -d该命令会自动按配置创建并启动所有容器。使用docker ps检查容器是否正常运行,也可以访问 http://localhost:5601/ 确认 Kibana 是否就绪。教程结束后,记得用以下命令停止并移除全部容器:
docker-compose down2. 下载连接器 JAR 并放置到 Flink lib 目录
Flink SQL Client 需要对应的连接器 JAR 才能解析sqlserver-cdc与elasticsearch-7两种 connector。将以下 JAR 包下载并放入<FLINK_HOME>/lib:
flink-sql-connector-elasticsearch7-3.0.1-1.17.jar:Elasticsearch 7.x 的 Flink SQL 连接器;flink-sql-connector-sqlserver-cdc:SqlServer CDC 连接器,从 Maven Central 仓库获取对应版本。
注意:上述下载链接仅对稳定发布版本(stable releases)有效;SNAPSHOT 版本需要基于本仓库的 master 或 release 分支自行构建。构建产物位于
flink-cdc-connect/flink-cdc-source-connectors/flink-sql-connector-sqlserver-cdc/模块下,该模块中的SqlServerTableFactory等源码正是被该 JAR 打包的核心实现(详见 SqlServerTableFactory.java)。
放置完成后重启 Flink 集群与 Flink SQL CLI,使新 JAR 生效。
3. 启用 SqlServer CDC 并准备测试数据
SQL Server 数据库默认不开启 CDC,需要执行系统存储过程启用。本教程在inventory数据库上开启数据库级 CDC,再对products与orders两张表开启表级 CDC:
-- Sqlserver CREATE DATABASE inventory; GO USE inventory; EXEC sys.sp_cdc_enable_db; -- Create and populate our products using a single insert with many rows CREATE TABLE products ( id INTEGER IDENTITY(101,1) NOT NULL PRIMARY KEY, name VARCHAR(255) NOT NULL, description VARCHAR(512), weight FLOAT ); INSERT INTO products(name,description,weight) VALUES ('scooter','Small 2-wheel scooter',3.14); INSERT INTO products(name,description,weight) VALUES ('car battery','12V car battery',8.1); INSERT INTO products(name,description,weight) VALUES ('12-pack drill bits','12-pack of drill bits with sizes ranging from #40 to #3',0.8); INSERT INTO products(name,description,weight) VALUES ('hammer','12oz carpenter''s hammer',0.75); INSERT INTO products(name,description,weight) VALUES ('hammer','14oz carpenter''s hammer',0.875); INSERT INTO products(name,description,weight) VALUES ('hammer','16oz carpenter''s hammer',1.0); INSERT INTO products(name,description,weight) VALUES ('rocks','box of assorted rocks',5.3); INSERT INTO products(name,description,weight) VALUES ('jacket','water resistent black wind breaker',0.1); INSERT INTO products(name,description,weight) VALUES ('spare tire','24 inch spare tire',22.2); EXEC sys.sp_cdc_enable_table @source_schema = 'dbo', @source_name = 'products', @role_name = NULL, @supports_net_changes = 0; -- Create some very simple orders CREATE TABLE orders ( id INTEGER IDENTITY(10001,1) NOT NULL PRIMARY KEY, order_date DATE NOT NULL, purchaser INTEGER NOT NULL, quantity INTEGER NOT NULL, product_id INTEGER NOT NULL, FOREIGN KEY (product_id) REFERENCES products(id) ); INSERT INTO orders(order_date,purchaser,quantity,product_id) VALUES ('16-JAN-2016', 1001, 1, 102); INSERT INTO orders(order_date,purchaser,quantity,product_id) VALUES ('17-JAN-2016', 1002, 2, 105); INSERT INTO orders(order_date,purchaser,quantity,product_id) VALUES ('19-FEB-2016', 1002, 2, 106); INSERT INTO orders(order_date,purchaser,quantity,product_id) VALUES ('21-FEB-2016', 1003, 1, 107); EXEC sys.sp_cdc_enable_table @source_schema = 'dbo', @source_name = 'orders', @role_name = NULL, @supports_net_changes = 0; GO要点解析:
EXEC sys.sp_cdc_enable_db开启数据库级 CDC;- 表级 CDC 通过
sys.sp_cdc_enable_table启用,参数@source_schema = 'dbo'、@source_name = 'products'指定被捕获的表;@role_name = NULL表示只有sysadmin或db_owner角色可访问变更表;@supports_net_changes = 0关闭净更改查询支持(本教程仅使用全部更改流); - 主键是 CDC 的硬性前提。
products与orders都定义了自增主键(IDENTITY),这也是 Flink CDC 增量快照分片(chunk)切分的依据。
从源码角度看,连接器在启动阶段会通过 SqlServerValidator.java 主动校验环境:它查询sys.databases判断目标库is_cdc_enabled是否为 1,未开启会抛出ValidationException("SqlServer database xxx do not enable cdc.");同时校验 SQL Server 主版本号,仅支持版本号大于等于 11 的实例。
4. 启动 Flink 集群并提交 Flink SQL 作业
4.1 开启 Checkpoint 并注册源表
启动 Flink 集群与 Flink SQL CLI 后,首先设置每 3 秒一次的 Checkpoint——这是 CDC 增量读取断点续传(exactly-once)的前提:
-- Flink SQL -- checkpoint every 3000 milliseconds Flink SQL> SET execution.checkpointing.interval = 3s; Flink SQL> CREATE TABLE products ( id INT, name STRING, description STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'sqlserver-cdc', 'hostname' = 'localhost', 'port' = '1433', 'username' = 'sa', 'password' = 'Password!', 'database-name' = 'inventory', 'table-name' = 'dbo.products' ); Flink SQL> CREATE TABLE orders ( id INT, order_date DATE, purchaser INT, quantity INT, product_id INT, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'sqlserver-cdc', 'hostname' = 'localhost', 'port' = '1433', 'username' = 'sa', 'password' = 'Password!', 'database-name' = 'inventory', 'table-name' = 'dbo.orders' );两个源表的 DDL 高度相似,核心差异仅在table-name:
'table-name' = 'dbo.products'采用<schema>.<table>的完整限定名格式,源码中 SqlServerTableFactory.java 将hostname/port/username/password/database-name/table-name声明为必填项(requiredOptions()),其余选项均有默认值;PRIMARY KEY (id) NOT ENFORCED必须与数据库主键一致,Flink CDC 依赖它进行增量快照的分片切分与变更事件的主键去重。
4.2 注册 Elasticsearch 结果表
Flink SQL> CREATE TABLE enriched_orders ( order_id INT, order_date DATE, purchaser INT, quantity INT, product_name STRING, product_description STRING, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'elasticsearch-7', 'hosts' = 'http://localhost:9200', 'index' = 'enriched_orders_1' );结果表enriched_orders对应 Elasticsearch 索引enriched_orders_1。由于 ES 的 upsert 语义要求主键字段参与索引文档_id的生成,这里声明了PRIMARY KEY (order_id) NOT ENFORCED。
4.3 提交流式 JOIN 作业
Flink SQL> INSERT INTO enriched_orders SELECT o.id,o.order_date,o.purchaser,o.quantity, p.name, p.description FROM orders AS o LEFT JOIN products AS p ON o.product_id = p.id;该语句构建了一个持续运行的流式作业:orders作为主表,products作为维表,通过product_id = id进行左连接,实时输出"订单 + 商品名 + 商品描述"的宽表结果并写入 Elasticsearch。Flink CDC 会将 SQL Server 事务日志中的每条变更(插入/更新/删除)转化为对应的流事件,Join 结果也会实时更新。
5. 在 Kibana 中验证同步结果
作业运行后,访问 http://localhost:5601/ 进入 Kibana:
- 创建索引模式(Index Pattern),匹配
enriched_orders_1; - 在 Discover 页面即可查看已同步的订单宽表数据,每行包含
order_id、order_date、purchaser、quantity、product_name、product_description字段。
初始快照阶段结束后,9 条商品数据与 4 条订单数据会通过 JOIN 展示为 4 条富化订单记录。
6. 修改源库数据,观察实时联动
保持 Flink 作业运行,回到 SQL Server 依次执行以下三类变更,Kibana 中的富化订单会在每一步操作后实时更新:
INSERT INTO orders(order_date,purchaser,quantity,product_id) VALUES ('22-FEB-2016', 1006, 22, 107); GO UPDATE orders SET quantity = 11 WHERE id = 10001; GO DELETE FROM orders WHERE id = 10004; GOINSERT:新增订单会立即出现在 Elasticsearch 索引中;UPDATE:订单10001的quantity从 1 变为 11,ES 中对应文档被 upsert 更新;DELETE:订单10004从 ES 中被删除。
这正是 CDC 流式同步的核心价值——Flink 作业一旦运行,源库的任何 DML 变更都会以近乎实时的延迟(取决于 Checkpoint 间隔与网络)传导到下游。仓库集成测试 SqlServerConnectorITCase.java 中testConsumingAllEvents用例完整覆盖了"快照读取 + 增量 UPDATE/INSERT/DELETE + 在线 Schema 变更"(ALTER TABLE增加列后重新开启 capture instance)的消费场景,验证了连接器对混合变更流的正确处理。
7. 连接器参数详解与源码印证
结合连接器官方文档 sqlserver-cdc.md 与 SqlServerTableFactory.java 的实现,常用参数如下:
| Option | 必填 | 默认值 | 类型 | 说明 |
|---|---|---|---|---|
connector | 是 | (无) | String | 固定为'sqlserver-cdc' |
hostname | 是 | (无) | String | SQL Server 的 IP 或主机名 |
username | 是 | (无) | String | 连接 SQL Server 的用户名 |
password | 是 | (无) | String | 连接密码 |
database-name | 是 | (无) | String | 要监控的数据库名 |
table-name | 是 | (无) | String | 要监控的表名,格式如"dbo.products" |
port | 否 | 1433 | Integer | SQL Server 端口号 |
server-time-zone | 否 | UTC | String | 数据库会话时区,如"Asia/Shanghai" |
scan.incremental.snapshot.enabled | 否 | true | Boolean | 是否启用并行快照(增量快照框架) |
chunk-meta.group.size | 否 | 1000 | Integer | 分片元数据分组大小,超限后元数据拆分为多组 |
chunk-key.even-distribution.factor.lower-bound | 否 | 0.05d | Double | 分片键均匀分布因子下界,用于判断表数据是否均匀分布以决定分片策略 |
chunk-key.even-distribution.factor.upper-bound | 否 | 1000.0d | Double | 分片键均匀分布因子上界,计算公式为(MAX(id) - MIN(id) + 1) / rowCount |
debezium.* | 否 | (无) | String | 透传 Debezium 属性给嵌入式引擎,如'debezium.snapshot.mode' = 'initial_only' |
scan.incremental.close-idle-reader.enabled | 否 | false | Boolean | 快照阶段结束后是否关闭空闲 reader;Flink ≥ 1.14 且开启execution.checkpointing.checkpoints-after-tasks-finish.enabled时生效 |
scan.incremental.snapshot.chunk.key-column | 否 | (无) | String | 快照分片键,默认取主键第一列;可用非主键列,但可能导致数据不一致 |
scan.incremental.snapshot.unbounded-chunk-first.enabled | 否 | true | Boolean | 快照阶段是否优先分配无界分片,可降低 TaskManager 处理最大无界分片时的 OOM 风险 |
scan.incremental.snapshot.backfill.skip | 否 | false | Boolean | 是否跳过快照阶段回填;跳过时快照期间的变更延后到日志阶段消费,仅保证 at-least-once |
scan.startup.mode | 否 | initial | String | 启动模式:initial/latest-offset/timestamp |
scan.startup.timestamp-millis | 否 | (无) | Long | 配合timestamp启动模式使用,指定起始时间戳(毫秒) |
从源码可以看到两点值得注意:
- 启动模式校验:
getStartupOptions()将字符串解析为StartupOptions.initial()/latest()/timestamp(),其中timestamp模式要求同时配置scan.startup.timestamp-millis,且只支持增量快照框架(即要求scan.incremental.snapshot.enabled = true),否则抛出校验异常; - 取值校验:
validateDistributionFactorUpper/Lower约束上界 ≥ 1.0、下界 ∈ [0, 1],validateIntegerOption保证 chunk 大小等整数选项大于下限,参数不合法时作业会在提交阶段直接失败而非运行后报错。
8. 启动模式与可用 Metadata
8.1 启动模式(Startup Reading Position)
scan.startup.mode决定连接器从何处开始读取:
initial(默认):先对捕获表做结构与数据的完整快照,再继续消费后续变更——即"存量 + 增量"全量读取;latest-offset:仅从当前时刻之后的变化开始读取,适用于只关心增量数据的场景;timestamp:从指定时间戳(scan.startup.timestamp-millis)之后开始读取。
注意:
scan.startup.mode底层依赖 Debezium 的snapshot.mode配置实现,两者不要同时使用,否则可能导致scan.startup.mode失效。集成测试testStartupFromLatestOffset(见 SqlServerConnectorITCase.java)验证了latest-offset模式下启动前插入的记录被丢弃、启动后的增量记录被完整消费。
8.2 元数据列(Metadata)
连接器支持以 VIRTUAL 只读列方式暴露以下元数据,定义方式如下:
CREATE TABLE products ( table_name STRING METADATA FROM 'table_name' VIRTUAL, schema_name STRING METADATA FROM 'schema_name' VIRTUAL, db_name STRING METADATA FROM 'database_name' VIRTUAL, operation_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL, id INT NOT NULL, name STRING, description STRING, weight DECIMAL(10,3) ) WITH ( 'connector' = 'sqlserver-cdc', 'hostname' = 'localhost', 'port' = '1433', 'username' = 'sa', 'password' = 'Password!', 'database-name' = 'inventory', 'table-name' = 'dbo.products' );可用元数据键(定义于 SqlServerReadableMetadata.java):
| Key | DataType | 说明 |
|---|---|---|
table_name | STRING NOT NULL | 行所属的表名 |
schema_name | STRING NOT NULL | 行所属的 schema 名 |
database_name | STRING NOT NULL | 行所属的数据库名 |
op_ts | TIMESTAMP_LTZ(3) NOT NULL | 数据库产生该变更的时间;若记录来自快照而非变更流,值恒为 0 |
源码实现中,这些元数据直接解析 DebeziumSourceRecord的source结构(TABLE_NAME_KEY/SCHEMA_NAME_KEY/DATABASE_NAME_KEY/TIMESTAMP_KEY)生成。
9. 数据类型映射
sqlserver-cdc连接器的 SQL Server 类型到 Flink SQL 类型映射如下(与集成测试testAllTypes中的full_types表 DDL 一一对应):
| SQLServer 类型 | Flink SQL 类型 |
|---|---|
char(n) | CHAR(n) |
varchar(n)/nvarchar(n)/nchar(n) | VARCHAR(n) |
text/ntext/xml | STRING |
decimal(p, s)/money/smallmoney | DECIMAL(p, s) |
numeric(p, s) | DECIMAL(p, s) |
float/real | DOUBLE |
bit | BOOLEAN |
int | INT |
tinyint/smallint | SMALLINT |
bigint | BIGINT |
date | DATE |
time(n) | TIME(n) |
datetime2/datetime/smalldatetime | TIMESTAMP(n) |
datetimeoffset | TIMESTAMP_LTZ(3) |
在定义 Flink 源表时,务必按此映射声明字段类型,避免精度丢失或类型不匹配导致的作业失败。
10. 关键使用限制与注意事项
- 快照期间无法进行 Checkpoint:当关闭增量快照框架(
scan.incremental.snapshot.enabled = false)时,快照阶段无可恢复位点,Checkpoint 会一直等待直至超时,超时 Checkpoint 默认会触发作业 Failover。若表数据量大,建议配置以下 Flink 参数规避:
execution.checkpointing.interval: 10min execution.checkpointing.tolerable-failed-checkpoints: 100 restart-strategy: fixed-delay restart-strategy.fixed-delay.attempts: 2147483647- 单线程读取限制:旧版 SourceFunction 形态的 SQL Server CDC 源无法并行读取(只有一个 task 能接收变更事件);基于增量快照框架的增量源(2.4.0 之后)支持设置并行度,例如
SqlServerSourceBuilder构建的源可setParallelism(2); - 无主键表:3.4.0 起支持无主键表,但必须配置
scan.incremental.snapshot.chunk.key-column指定一个非空字段。若该字段发生 UPDATE,只能保证 at-least-once 语义,建议下游声明主键并做幂等处理;对已有主键的表使用非主键列作为分片键可能造成数据不一致(同一行被两个分片以不同快照值输出,最终结果取决于处理顺序)。
11. 进一步扩展:DataStream API 形态
除 Flink SQL 外,连接器同样可作为 DataStream Source 使用。增量快照框架版本(2.4.0 之后)的构建方式如下(完整示例见 SqlServerSourceBuilder.java):
import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.cdc.connectors.base.options.StartupOptions; import org.apache.flink.cdc.connectors.sqlserver.source.SqlServerSourceBuilder; import org.apache.flink.cdc.connectors.sqlserver.source.SqlServerSourceBuilder.SqlServerIncrementalSource; import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema; public class SqlServerIncrementalSourceExample { public static void main(String[] args) throws Exception { SqlServerIncrementalSource<String> sqlServerSource = new SqlServerSourceBuilder() .hostname("localhost") .port(1433) .databaseList("inventory") .tableList("dbo.products") .username("sa") .password("Password!") .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.initial()) .build(); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // enable checkpoint env.enableCheckpointing(3000); // set the source parallelism to 2 env.fromSource( sqlServerSource, WatermarkStrategy.noWatermarks(), "SqlServerIncrementalSource") .setParallelism(2) .print() .setParallelism(1); env.execute("Print SqlServer Snapshot + Change Stream"); } }SqlServerSourceBuilder通过链式 API 组装hostname/port/databaseList/tableList/username/password/deserializer/startupOptions等配置,内部由SqlServerSourceConfigFactory生成配置,再经LsnFactory与SqlServerDialect构建出基于 LSN(Log Sequence Number)断点管理的增量源。
12. 清理环境
实验完成后,按顺序清理:
- 在 Flink SQL CLI 中取消(Cancel)流式 JOIN 作业;
- 停止并移除容器:
docker-compose down; - 如需彻底删除 Docker 卷中的数据,可追加
-v参数。
至此,你已经完整走通了"SQL Server → Flink CDC → Elasticsearch → Kibana"的实时数据管道:掌握了数据库/表级 CDC 的启用、Flink SQL 中sqlserver-cdc源表与elasticsearch-7结果表的定义、流式维表关联的编写方法,以及增量 DML 变更在下游的实时联动验证。更多参数细节可继续阅读连接器完整文档 sqlserver-cdc.md 与对应源码模块 flink-connector-sqlserver-cdc。
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考