SQLServer CDC 连接器可以读取 SQLServer 数据库的快照数据和增量数据。本文为您介绍如何使用SQLServer CDC连接器。
背景信息
SQLServer CDC连接器支持的信息如下。
类别 | 详情 |
支持类型 | 源表 |
运行模式 | 仅支持流模式 |
数据格式 | 暂不适用 |
特有监控指标 | |
API种类 | SQL、数据摄入 |
是否支持更新或删除结果表数据 | 不涉及 |
前提条件
SQLServer CDC连接器使用前,需要捕获的数据库和数据表需要启用变更数据捕获功能(CDC)。
数据表开启CDC功能可以执行如下SQL:
RDS SQLServer
EXEC sp_rds_cdc_enable_db; -- 对应数据库开启CDC
EXEC sys.sp_cdc_enable_table
@source_schema = N'dbo', -- 指定源表的 schema
@source_name = N'MyTable', -- 指定要捕获的表名
@role_name = N'MyRole', -- 指定角色 MyRole
@supports_net_changes = 1;自建SQLServer
EXEC sys.sp_cdc_enable_table
@source_schema = N'dbo', -- 指定源表的 schema
@source_name = N'MyTable', -- 指定要捕获的表名
@role_name = N'MyRole', -- 指定角色 MyRole,可将需要对源表捕获列
-- 拥有 SELECT 权限的用户添加到该角色。
-- sysadmin 或 db_owner 角色的用户也可以
-- 访问指定的变更表。设置为 NULL 则仅允许
-- sysadmin 或 db_owner 完全访问捕获信息。
@filegroup_name = N'MyDB_CT', -- 指定文件组,SQL Server 会将变更表
-- 放置在该文件组中。该文件组必须已存在。
-- 建议不要将变更表与源表放在同一文件组中。
@supports_net_changes = 0;验证CDC表的访问权限可使用如下SQL:
EXEC sys.sp_cdc_help_change_data_capture;该查询返回数据库中已启用 CDC 的每个表的配置信息。如果对应表不存在,请检查表是否启用CDC,且用户是否拥有访问捕获实例和CDC表的权限。
使用限制
使用前需确保目标数据库已启用 CDC。
SQL功能从VVR 11.7版本开始支持,数据摄入功能从VVR 11.8版本开始支持。
由于连接器依赖 SQL Server 的 Change Data Capture 功能,存在以下限制:
Standard 版时需要 SQL Server 2016 SP1 及以上。
Enterprise 版需要 SQL Server 2012 及以上。
SQL
语法结构
CREATE TABLE sqlserver_cdc_source (
id INT,
order_date DATE,
purchaser INT,
quantity INT,
product_id INT,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'sqlserver-cdc',
'hostname' = '<hostname>',
'port' = '<port>',
'username' = '<username>',
'password' = '<password>',
'database-name' = '<database name>',
'table-name' = 'dbo.orders'
);WITH参数
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 | ||||||
connector | connector类型。 | STRING | 是 | 无 | 固定值为 | ||||||
hostname | SQLServer数据库的IP地址或者Hostname。 | STRING | 是 | 无 | 无。 | ||||||
username | SQLServer数据库服务的用户名。 | STRING | 是 | 无 | 无。 | ||||||
password | SQLServer数据库服务的密码。 | STRING | 是 | 无 | 无。 | ||||||
database-name | 数据库名称。 | STRING | 是 | 无 | 数据库名称。 | ||||||
table-name | 需要捕获的表名。 | STRING | 是 | 无 | 多个表匹配可用
| ||||||
scan.startup.mode | 消费数据时的启动模式。 | STRING | 否 | initial | 参数取值如下:
重要
| ||||||
scan.startup.timestamp-millis | 启动时间戳(毫秒)。 | LONG | 否 | 无 | 仅在 scan.startup.mode 为 timestamp 时需要配置。表示自 epoch(1970-01-01 00:00:00 UTC)以来的毫秒数。 | ||||||
port | SQLServer数据库服务的端口号。 | INTEGER | 否 | 1433 | 无。 | ||||||
server-time-zone | 数据库服务器的会话时区,如 | STRING | 否 | 无 | 无。 | ||||||
scan.incremental.snapshot.enabled | 是否开启增量快照。 | BOOLEAN | 否 | true | 参数取值如下:
| ||||||
scan.incremental.snapshot.chunk.size | 表快照的 chunk 大小(行数)。 | INTEGER | 否 | 8096 | 读取快照时,表会被拆分为多个 chunk。 | ||||||
scan.snapshot.fetch.size | 每次轮询的最大拉取行数。 | INTEGER | 否 | 1024 | 无。 | ||||||
chunk-meta.group.size | chunk 元数据的分组大小,超过该大小时元数据会被分成多组。 | INTEGER | 否 | 1000 | 无。 | ||||||
chunk-key.even-distribution.factor.lower-bound | chunk key 分布因子下界。 | DOUBLE | 否 | 0.05 | 分布因子用于判断表数据是否均匀分布。均匀分布时使用均匀计算优化,不均匀时会触发拆分查询。计算公式: | ||||||
chunk-key.even-distribution.factor.upper-bound | chunk key 分布因子上界。 | DOUBLE | 否 | 1000.0 | 分布因子用于判断表数据是否均匀分布。均匀分布时使用均匀计算优化,不均匀时会触发拆分查询。计算公式: | ||||||
scan.incremental.close-idle-reader.enabled | 是否在快照结束后关闭空闲的Reader。 | BOOLEAN | 否 | false | 该配置生效需要设置 | ||||||
scan.incremental.snapshot.chunk.key-column | 指定某一列作为快照阶段切分分片的切分列。 | STRING | 否 | 无 | 表快照的 chunk key 列。读取快照时按 chunk key 将表拆分为多个 chunk。默认使用主键的第一列,必须从主键列中选择。 | ||||||
scan.incremental.snapshot.unbounded-chunk-first.enabled | 是否在快照阶段优先分配无界 chunk,有助于降低 TaskManager 在读取最大无界 chunk 时发生 OOM 的风险。 | BOOLEAN | 否 | true | 无。 | ||||||
scan.incremental.snapshot.backfill.skip | 是否跳过全量阶段的日志读取。 | BOOLEAN | 否 | false | 参数取值如下:
| ||||||
connect.timeout | 连接超时时间。 | DURATION | 否 | 30s | 无。 | ||||||
connect.max-retries | 连接最大重试次数。 | INTEGER | 否 | 3 | 无。 | ||||||
connection.pool.size | 连接池大小。 | INTEGER | 否 | 20 | 无。 | ||||||
debezium.* | Debezium属性参数。 | STRING | 否 | 无 | 更细粒度控制Debezium客户端的行为,详情请参见配置属性。 无特别需求不建议自行配置Debezium参数,可能导致无法正常读取数据。 |
类型映射
SQLServer和Flink字段类型对应关系如下。
SQLServer字段类型 | Flink字段类型 |
bit | BOOLEAN |
tinyint | SMALLINT |
smallint | |
int | INT |
bigint | BIGINT |
float | DOUBLE |
real | |
numeric | DECIMAL(p,s) |
decimal(p,s) | |
money | |
smallmoney | |
date | DATE |
time(n) | TIME(n) |
datetime2 | TIMESTAMP(n) |
datetime | |
smalldatetime | |
datetimeoffset | TIMESTAMP_LTZ(3) |
char(n) | CHAR(n) |
varchar(n) | VARCHAR(n) |
nvarchar(n) | |
nchar(n) | |
text | STRING |
ntext | |
xml |
元数据
SQLServer CDC 源表支持以下元数据列,在建表 DDL 中通过 METADATA FROM 语法声明。
Key | 数据类型 | 说明 |
database_name | STRING NOT NULL | 该行数据所属的数据库名称。 |
schema_name | STRING NOT NULL | 该行数据所属的 Schema 名称。 |
table_name | STRING NOT NULL | 该行数据所属的表名称。 |
op_ts | TIMESTAMP_LTZ(3) NOT NULL | 该行数据在数据库中的变更时间。如果该行数据来自快照(全量阶段),则值为 0(即 1970-01-01 00:00:00)。 |
op_type | STRING NOT NULL | 事件的类型,共计4种:+I(INSERT)、-D(DELETE)、-U(UPDATE_BEFORE)、+U(UPDATE_AFTER)。 |
Flink SQL 中声明SQLServer 元数据列示例:
CREATE TEMPORARY TABLE sqlserver_cdc_source (
id INT,
name STRING,
-- 元数据列
db_name STRING METADATA FROM 'database_name' VIRTUAL,
tbl_name STRING METADATA FROM 'table_name' VIRTUAL,
op_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL,
op_type STRING METADATA FROM 'op_type' VIRTUAL,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'sqlserver-cdc',
'hostname' = '<yourHostname>',
'port' = '1433',
'username' = '<yourUsername>',
'password' = '<yourPassword>',
'database-name' = '<yourDatabaseName>',
'schema-name' = 'dbo',
'table-name' = '<yourTableName>'
);数据摄入
语法结构
source:
type: sqlserver
name: SQLServer Source
hostname: 127.0.0.1
port: 1433
username: sa
password: Password!
tables: inventory.dbo.orders
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://127.0.0.1:9030
load-url: 127.0.0.1:8030
username: root
password: pass说明:tables 配置项使用三段式命名 database.schema.table,支持正则表达式匹配多表。例如 inventory.dbo.\.* 表示捕获 inventory 数据库 dbo schema 下的所有表。
配置项
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
type | Source 类型。 | STRING | 是 | 无 | 固定值为 sqlserver。 |
hostname | SQLServer 数据库的 IP 地址或者 Hostname。 | STRING | 是 | 无 | 无。 |
username | SQLServer 数据库服务的用户名。 | STRING | 是 | 无 | 无。 |
password | SQLServer 数据库服务的密码。 | STRING | 是 | 无 | 无。 |
tables | 需要捕获的表名。 | STRING | 是 | 无 | 使用三段式命名 database.schema.table,其中 database 必须为固定名称,不支持正则,schema和table支持正则表达式,。多个表用逗号(,)分隔。例如 |
port | SQLServer 数据库服务的端口号。 | INTEGER | 否 | 1433 | 无。 |
tables.exclude | 需要排除的表名。 | STRING | 否 | 无 | 语法与 tables 相同,用于从已匹配的表中排除指定表。 |
scan.startup.mode | 消费数据时的启动模式。 | STRING | 否 | initial | 参数取值如下:
|
scan.startup.timestamp-millis | 启动时间戳(毫秒)。 | LONG | 否 | 无 | 仅在 scan.startup.mode 为 timestamp 时需要配置。 |
server-time-zone | 数据库服务器的会话时区。 | STRING | 否 | 无 | 例如 "Asia/Shanghai"。如果未设置,则使用系统默认时区。 |
scan.incremental.snapshot.chunk.size | 表快照的 chunk 大小(行数)。 | INTEGER | 否 | 8096 | 读取快照时,表会被拆分为多个 chunk。 |
scan.snapshot.fetch.size | 每次轮询的最大拉取行数。 | INTEGER | 否 | 1024 | 无。 |
scan.incremental.snapshot.chunk.key-column | 指定快照阶段的切分列。 | STRING | 否 | 无 | 默认使用主键的第一列,必须从主键列中选择。 |
scan.incremental.snapshot.backfill.skip | 是否跳过全量阶段的日志回填。 | BOOLEAN | 否 | true | 设置为 true 时,全量阶段不读取日志,仅保证 At-Least-Once 语义。 |
schema-change.enabled | 是否发送 Schema 变更事件。 | BOOLEAN | 否 | true | 设置为 false 时,上游的 DDL 变更不会同步到下游。 |
connect.timeout | 连接超时时间。 | DURATION | 否 | 30s | 无。 |
connect.max-retries | 连接最大重试次数。 | INTEGER | 否 | 3 | 无。 |
connection.pool.size | 连接池大小。 | INTEGER | 否 | 20 | 无。 |
scan.incremental.close-idle-reader.enabled | 是否在快照结束后关闭空闲的 Reader。 | BOOLEAN | 否 | false | 该配置生效需要设置 execution.checkpointing.checkpoints-after-tasks-finish.enabled 为 true。 |
scan.incremental.snapshot.unbounded-chunk-first.enabled | 是否在快照阶段优先分配无界 chunk。 | BOOLEAN | 否 | false | 有助于降低 TaskManager 在读取最大无界 chunk 时发生 OOM 的风险。 |
scan.newly-added-table.enabled | 是否扫描新增的表。 | BOOLEAN | 否 | false | 仅在从 Savepoint/Checkpoint 恢复时有效。 |
metadata.list | 需要传递给下游的元数据列表。 | STRING | 否 | 无 | 多个元数据用逗号分隔。可选值:database_name、schema_name、table_name、op_ts。 |
chunk-meta.group.size | chunk 元数据的分组大小。 | INTEGER | 否 | 1000 | 无。 |
chunk-key.even-distribution.factor.upper-bound | chunk key 分布因子上界。 | DOUBLE | 否 | 1000.0 | 分布因子用于判断表数据是否均匀分布。计算公式:(MAX(id) - MIN(id) + 1) / rowCount。 |
chunk-key.even-distribution.factor.lower-bound | chunk key 分布因子下界。 | DOUBLE | 否 | 0.05 | 同上。 |
debezium.* | Debezium 属性参数。 | STRING | 否 | 无 | 更细粒度控制 Debezium 客户端的行为。无特别需求不建议自行配置。 |
jdbc.properties.* | JDBC 连接属性。 | STRING | 否 | 无 | 传递给 JDBC 驱动的额外属性。 |
类型映射
SQLServer 和 Flink CDC 字段类型对应关系如下。
SQLServer 字段类型 | Flink CDC 字段类型 |
bit | BOOLEAN |
tinyint | SMALLINT |
smallint | SMALLINT |
int | INT |
bigint | BIGINT |
real | FLOAT |
float | DOUBLE |
numeric(p,s) | DECIMAL(p,s) |
decimal(p,s) | DECIMAL(p,s) |
money | DECIMAL(19,4) |
smallmoney | DECIMAL(10,4) |
date | DATE |
time(n) | TIME(n) |
datetime2(n) | TIMESTAMP(n) |
datetime | TIMESTAMP(3) |
smalldatetime | TIMESTAMP(0) |
datetimeoffset(n) | TIMESTAMP_LTZ(n) |
char(n) | CHAR(n) |
nchar(n) | CHAR(n) |
varchar(n) | VARCHAR(n) |
nvarchar(n) | VARCHAR(n) |
text | STRING |
ntext | STRING |
xml | STRING |
uniqueidentifier | STRING |
binary(n) | BYTES |
varbinary(n) | BYTES |
image | BYTES |
timestamp / rowversion | BYTES |
元数据
在Flink CDC数据摄入作业中,元数据分为两类:
1. Flink CDC 框架支持的元数据列:由 Flink CDC框架根据每条事件的属性自动推断而来,不需要在source模块配置 metadata.list, 参考 Flink CDC的 元数据列(Metadata Column)。
2. 数据源支持的元数据列:从数据源端变更事件中读取,需要在source模块配置 metadata.list,不同数据源支持的元数据不同,SQLServer 支持的元数据如下,声明多个元数据时用逗号分隔。
Key | 数据类型 | 说明 |
database_name | STRING NOT NULL | 该行数据所属的数据库名称。 |
schema_name | STRING NOT NULL | 该行数据所属的 Schema 名称。 |
table_name | STRING NOT NULL | 该行数据所属的表名称。 |
op_ts | TIMESTAMP_LTZ(3) NOT NULL | 该行数据在数据库中的变更时间。如果该行数据来自快照(全量阶段),则值为 0(即 1970-01-01 00:00:00)。 |
Flink CDC 声明SQLServer元数据示例:
source:
type: sqlserver
name: SQLServer Source
hostname: 127.0.0.1
port: 1433
username: sa
password: Password!
tables: inventory.dbo.orders
-- SQLServer 数据源支持的元数据,声明多个时用逗号分隔
metadata.list: database_name,schema_name,table_name,op_ts
transform:
- source: inventory.dbo.orders
-- * 表示引用源表全部字段,table_name,op_ts表示引用source提供的元数据列,op_type 表示引用框架提供的元数据列
projection: "*, table_name AS table_name, op_ts AS op_ts, __data_event_type__ AS op_type"
说明:__data_event_type__ 是Flink CDC框架支持的元数据列,可以用于判断事件类型,Flink CDC 有三种事件类型:插入,更新,删除。其中更新事件总是在一条记录中,这与Flink SQL的更新保存为两条独立记录不同,优势是可以保留完整的更新语义,使得Flink CDC支持将原始更新事件同步到下游系统。