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 |
参数取值如下:
backfill仅在单个分片(chunk)快照查询期间生效,不覆盖整个全量读取过程。跳过backfill后,分片快照SQL执行时读到该时刻表的最新数据;分片已读完之后该分片上发生的更新,不再在全量阶段合并,会在进入增量阶段后从变更日志中读取。例如,chunk5快照期间发生的更新会直接体现在chunk5的最新数据中;若已读到chunk80时chunk5才发生更新,该更新会在增量阶段通过变更日志补回。 开启后,分片扫描期间及之后的变更在增量阶段仍会通过变更日志下发,可能与快照数据重复,仅提供at-least-once语义。请确认下游支持按主键幂等写入后再开启。 backfill仅在单个分片(chunk)快照查询期间生效,不覆盖整个全量读取过程。跳过backfill后,分片快照SQL执行时读到该时刻表的最新数据;分片已读完之后该分片上发生的更新,不再在全量阶段合并,会在进入增量阶段后从变更日志中读取。例如,chunk5快照期间发生的更新会直接体现在chunk5的最新数据中;若已读到chunk80时chunk5才发生更新,该更新会在增量阶段通过变更日志补回。 开启后,分片扫描期间及之后的变更在增量阶段仍会通过变更日志下发,可能与快照数据重复,仅提供at-least-once语义。请确认下游支持按主键幂等写入后再开启。 |
||||||
|
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 |
|
real |
FLOAT |
|
float |
DOUBLE |
|
numeric(p,s) |
DECIMAL(p,s) |
|
decimal(p,s) |
|
|
money |
DECIMAL(19,4) |
|
smallmoney |
DECIMAL(10,4) |
|
date |
DATE |
|
time(p) |
TIME(p) |
|
datetime2(p) |
TIMESTAMP(p) |
|
datetime |
TIMESTAMP(3) |
|
smalldatetime |
TIMESTAMP(0) |
|
datetimeoffset(p) |
TIMESTAMP_LTZ(p) |
|
char(n) |
CHAR(n) |
|
nchar(n) |
|
|
varchar(n) |
VARCHAR(n) |
|
nvarchar(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支持将原始更新事件同步到下游系统。