MySQL连接器作为数据源可以在数据摄入YAML作业中使用。
前提条件
在使用MySQL CDC源表前,必须先按照配置MySQL进行操作,这些操作主要为了满足使用MySQL CDC源表的前提条件
RDS MySQL
与实时计算Flink版进行网络探测,确保网络连通。
MySQL版本要求:5.6,5.7,8.0.x,8.4。
需开启Binlog(默认开启)。
Binlog格式需要为ROW(默认)。
设置binlog_row_image为FULL(默认)。
关闭Binary Log Transaction Compression。(8.0.20及以上引入,默认关闭)。
已创建MySQL用户,并授予了SELECT、SHOW DATABASES、REPLICATION SLAVE和REPLICATION CLIENT权限。
已创建MySQL数据库和表,详情请参见RDS MySQL创建数据库和账号。(请使用高权限账号来创建MySQL数据库,避免因权限不足而导致操作失败。)
已设置IP白名单,详情请参见RDS MySQL白名单设置。
PolarDB MySQL
与实时计算Flink版进行网络探测,确保网络连通。
MySQL版本要求:5.6,5.7,8.0.x,8.4。
需开启Binlog(默认关闭)。
Binlog格式需要为ROW(默认)。
设置binlog_row_image为FULL(默认)。
关闭Binary Log Transaction Compression。(8.0.20及以上引入,默认关闭)。
已创建MySQL用户,并授予了SELECT、SHOW DATABASES、REPLICATION SLAVE和REPLICATION CLIENT权限。
已创建MySQL数据库和表,详情请参见PolarDB MySQL创建数据库和账号。(请使用高权限账号来创建MySQL数据库,避免因权限不足而导致操作失败。)
已设置IP白名单,详情请参见PolarDB MySQL白名单设置。
自建MySQL
与实时计算Flink版进行网络探测,确保网络连通。
MySQL版本要求:5.6,5.7,8.0.x,8.4。
需开启Binlog(默认关闭)。
Binlog格式需要为ROW(默认为STATEMENT)。
设置binlog_row_image为FULL(默认)。
关闭Binary Log Transaction Compression。(8.0.20及以上引入,默认关闭)。
已创建MySQL用户,并授予了SELECT、SHOW DATABASES、REPLICATION SLAVE和REPLICATION CLIENT权限。
已创建MySQL数据库和表。(请使用高权限账号来创建MySQL数据库,避免因权限不足而导致操作失败。)
已设置IP白名单,详情请参见自建MySQL白名单设置。
使用限制
通用限制
MySQL CDC连接器目前暂不支持Binary Log Transaction Compression(二进制日志事务压缩) 功能。因此,在使用MySQL CDC连接器消费增量数据时,请务必确保已关闭Binary Log Transaction Compression配置,否则可能导致增量数据无法正常获取。
RDS MySQL的限制
对于RDS MySQL,不建议通过备库或只读从库读取数据。因为RDS MySQL的备库和只读从库Binlog保留时间默认很短,可能由于Binlog过期清理,导致作业无法消费Binlog数据而报错。
RDS MySQL默认开启了主从并行同步功能,且不保证主从事务顺序一致,可能导致主从切换后并Checkpoint恢复时漏读部分数据。您可以手动打开RDS MySQL的slave_preserve_commit_order选项来规避此问题。
PolarDB MySQL的限制
MySQL CDC源表不支持读取PolarDB MySQL版1.0.19及以前版本的多主架构集群(什么是多主集群?)。PolarDB MySQL版1.0.19及更早版本的多主架构集群产生的Binlog可能出现重复Table ID,导致CDC源表Schema映射错误,从而解析Binlog数据报错。
开源MySQL的限制
在默认配置下,MySQL进行主从Binlog复制时,总是保持Transaction顺序。若MySQL副本启用了并行复制(slave_parallel_workers> 1)但未开启 slave_preserve_commit_order=ON,其事务提交顺序可能与主库不一致。Flink CDC 从检查点恢复时会因顺序错乱而漏读数据。推荐在MySQL副本上设置 slave_preserve_commit_order = ON。或设置 slave_parallel_workers = 1(会牺牲复制性能)。
注意事项
在全量阶段时,不支持保存savepoint后,在源表增加新表或删除表再从savepoint重启的操作,这会导致作业无法正常读取数据。
数据摄入
MySQL连接器作为数据源可以在数据摄入YAML作业中使用。
语法结构
source:
type: mysql
name: MySQL Source
hostname: localhost
port: 3306
username: <username>
password: <password>
tables: adb.\.*, bdb.user_table_[0-9]+, [app|web].order_\.*
server-id: 5401-5404
sink:
type: xxx配置项
参数 | 说明 | 是否必填 | 数据类型 | 默认值 | 备注 |
type | 数据源类型。 | 是 | STRING | 无 | 固定值为mysql。 |
name | 数据源名称。 | 否 | STRING | 无 | 无。 |
hostname | MySQL数据库的IP地址或者Hostname。 | 是 | STRING | 无 | 建议填写专有网络VPC地址。 说明 如果MySQL与实时Flink版不在同一VPC,需要先打通跨VPC的网络或者使用公网的形式访问,详情请参见空间管理与操作和Flink全托管集群如何访问公网?。 |
username | MySQL数据库服务的用户名。 | 是 | STRING | 无 | 无。 |
password | MySQL数据库服务的密码。 | 是 | STRING | 无 | 无。 |
tables | 需要同步的MySQL数据表。 | 是 | STRING | 无 |
说明
|
tables.exclude | 需要在同步的表中排除的表。 | 否 | STRING | 无 |
说明 点号用于分割数据库名和表名,如果需要用点号匹配任意字符,需要对点号使用反斜杠进行转译。如:db0.\.*, db1.user_table_[0-9]+, db[1-2].[app|web]order_\.*。 |
port | MySQL数据库服务的端口号。 | 否 | INTEGER | 3306 | 无。 |
schema-change.enabled | 是否发送Schame变更事件。 | 否 | BOOLEAN | true | 无。 |
server-id | 数据库客户端的用于同步的数字ID或范围。 | 否 | STRING | 默认会随机生成一个5400~6400的值。 | 该ID必须是MySQL集群中全局唯一的。建议针对同一个数据库的每个作业都设置一个不同的ID。该参数也支持ID范围的格式,例如5400-5408。 说明 在开启增量读取模式时支持多并发读取,此时推荐设定为ID范围,使得每个并发使用不同的ID。 |
jdbc.properties.* | JDBC URL中的自定义连接参数。 | 否 | STRING | 无 | 您可以传递自定义的连接参数,例如不使用SSL协议,则可配置为'jdbc.properties.useSSL' = 'false'。 支持的连接参数请参见MySQL Configuration Properties。 |
debezium.* | Debezium读取Binlog的自定义参数。 | 否 | STRING | 无 | 您可以传递自定义的Debezium参数,例如使用'debezium.event.deserialization.failure.handling.mode'='ignore'来指定解析错误时的处理逻辑。 警告 请不要随意修改debeizum参数,这可能导致连接器读取数据错误。例如,debezium.binlog.buffer.size参数是禁止配置的。 |
scan.incremental.snapshot.chunk.size | 每个chunk的大小(包含的行数)。 | 否 | INTEGER | 8096 | MySQL表会被切分成多个chunk读取。在读完chunk的数据之前,chunk的数据会先缓存在内存中。 每个chunk包含的行数越少,则表中的chunk的总数量越大,尽管这会降低故障恢复的粒度,但可能导致内存OOM和整体的吞吐量降低。因此,您需要进行权衡,并设置合理的chunk大小。 |
scan.snapshot.fetch.size | 当读取表的全量数据时,每次最多拉取的记录数。 | 否 | INTEGER | 1024 | 无。 |
scan.startup.mode | 消费数据时的启动模式。 | 否 | STRING | initial | 参数取值如下:
重要 对于earliest-offset,specific-offset和timestamp启动模式,如果启动时刻和指定的启动位点时刻的表结构不同,作业会因为表结构不同而报错。换一句话说,使用这三种启动模式,需要保证在指定的Binlog消费位置到作业启动的时间之间,对应表不能发生表结构变更。 |
scan.startup.specific-offset.file | 使用指定位点模式启动时,启动位点的Binlog文件名。 | 否 | STRING | 无 | 使用该配置时,scan.startup.mode必须配置为specific-offset。文件名格式例如 |
scan.startup.specific-offset.pos | 使用指定位点模式启动时,启动位点在指定Binlog文件中的偏移量。 | 否 | INTEGER | 无 | 使用该配置时,scan.startup.mode必须配置为specific-offset。 |
scan.startup.specific-offset.gtid-set | 使用指定位点模式启动时,启动位点的GTID集合。 | 否 | STRING | 无 | 使用该配置时,scan.startup.mode必须配置为specific-offset。GTID集合格式例如 |
scan.startup.timestamp-millis | 使用指定时间模式启动时,启动位点的毫秒时间戳。 | 否 | LONG | 无 | 使用该配置时,scan.startup.mode必须配置为timestamp。时间戳单位为毫秒。 重要 在使用指定时间时,MySQL CDC会尝试读取每个Binlog文件的初始事件以确定其时间戳,最终定位至指定时间对应的Binlog文件。请保证指定的时间戳对应的Binlog文件在数据库中没有被清理且可以被读取到。 |
server-time-zone | 数据库在使用的会话时区。 | 否 | STRING | 如果您没有指定该参数,则系统默认使用Flink作业运行时的环境时区作为数据库服务器时区,即您选择的可用区所在的时区。 | 例如Asia/Shanghai,该参数控制了MySQL中的TIMESTAMP类型如何转成STRING类型。 |
scan.startup.specific-offset.skip-events | 从指定的位点读取时,跳过多少Binlog事件。 | 否 | INTEGER | 无 | 使用该配置时,scan.startup.mode必须配置为specific-offset。 |
scan.startup.specific-offset.skip-rows | 从指定的位点读取时,跳过多少行变更(一个Binlog事件可能对应多行变更)。 | 否 | INTEGER | 无 | 使用该配置时,scan.startup.mode必须配置为specific-offset。 |
connect.timeout | 连接MySQL数据库服务器超时时,重试连接之前等待超时的最长时间。 | 否 | DURATION | 30s | 无。 |
connect.max-retries | 连接MySQL数据库服务时,连接失败后重试的最大次数。 | 否 | INTEGER | 3 | 无。 |
connection.pool.size | 数据库连接池大小。 | 否 | INTEGER | 20 | 数据库连接池用于复用连接,可以降低数据库连接数量。 |
heartbeat.interval | Source通过心跳事件推动Binlog位点前进的时间间隔。 | 否 | DURATION | 30s | 心跳事件用于推动Source中的Binlog位点前进,这对MySQL中更新缓慢的表非常有用。对于更新缓慢的表,Binlog位点无法自动前进,通过够心跳事件可以推到Binlog位点前进,可以避免Binlog位点不前进引起Binlog位点过期问题,Binlog位点过期会导致作业失败无法恢复,只能无状态重启。 |
rds.region-id | 阿里云RDS MySQL实例所在的地域ID。 | 使用读取OSS归档日志功能时必填。 | STRING | 无 | 地域ID请参见地域和可用区。 重要 因为MySQL CDC的GTID字符串是随机生成,不像binlog文件位点是单调递增,在定位一个GTID所在文件时需要下载并解析全量oss归档日志,其开销和耗时非常大,依赖GTID位点的功能不具备可行性。所以OSS归档日志功能仅支持指定时间戳和指定binlog文件位点启动,不支持指定GTID启动,也不支持归档日志中存在过主从切换的场景,因为MySQL主从切换依赖GTID,请您在使用该功能前谨慎评估。 |
rds.access-key-id | 阿里云RDS MySQL账号Access Key ID。 | 使用读取OSS归档日志功能时必填。 | STRING | 无 | 详情请参见如何查看AccessKey ID和AccessKey Secret信息? 重要 为了避免您的AK信息泄露,建议您通过密钥管理的方式填写AccessKey ID取值,详情请参见变量管理。 |
rds.access-key-secret | 阿里云RDS MySQL账号Access Key Secret。 | 使用读取OSS归档日志功能时必填。 | STRING | 无 | 详情请参见如何查看AccessKey ID和AccessKey Secret信息? 重要 为了避免您的AK信息泄露,建议您通过密钥管理的方式填写AccessKey Secret取值,详情请参见变量管理。 |
rds.db-instance-id | 阿里云RDS MySQL实例ID。 | 使用读取OSS归档日志功能时必填。 | STRING | 无 | 无。 |
rds.main-db-id | 阿里云RDS MySQL实例主库编号。 | 否 | STRING | 无 | 获取主库编号详情请参见RDS MySQL日志备份。 说明 如果未填写,VVR 11.7及以上版本会根据RDS MySQL连接信息自动查询主库编号。 |
rds.download.timeout | 从OSS下载单个归档日志的超时时间。 | 否 | DURATION | 60s | 无。 |
rds.endpoint | 获取OSS Binlog信息的服务接入点。 | 否 | STRING | 无 | 可选值详情请参见服务接入点。 |
rds.binlog-directory-prefix | 保存Binlog文件的目录前缀。 | 否 | STRING | rds-binlog- | 无。 |
rds.use-intranet-link | 是否使用内网下载Binlog文件。 | 否 | BOOLEAN | true | 无。 |
rds.binlog-directories-parent-path | 保存Binlog文件的父目录的绝对路径。 | 否 | STRING | 无 | 无。 |
chunk-meta.group.size | chunk元信息的大小。 | 否 | INTEGER | 1000 | 如果元信息大于该值,元信息会分为多份传递。 |
chunk-key.even-distribution.factor.lower-bound | 是否可以均匀分片的chunk分布因子的下限。 | 否 | DOUBLE | 0.05 | 分布因子小于该值会使用非均匀分片。 chunk分布因子 = (MAX(chunk-key) - MIN(chunk-key) + 1) / 总数据行数。 |
chunk-key.even-distribution.factor.upper-bound | 是否可以均匀分片的chunk分布因子的上限。 | 否 | DOUBLE | 1000.0 | 分布因子大于该值会使用非均匀分片。 chunk分布因子 = (MAX(chunk-key) - MIN(chunk-key) + 1) / 总数据行数。 |
scan.incremental.close-idle-reader.enabled | 是否在快照结束后关闭空闲的Reader。 | 否 | BOOLEAN | false | 该配置生效,需要设置 |
scan.only.deserialize.captured.tables.changelog.enabled | 在增量阶段,是否仅对指定表的变更事件进行反序列化。 | 否 | BOOLEAN |
| 参数取值如下:
|
scan.parallel-deserialize-changelog.enabled | 在增量阶段,是否使用多线程对变更事件进行解析。 | 否 | BOOLEAN | false | 参数取值如下:
说明 仅Flink计算引擎VVR 8.0.11及以上版本支持。 |
scan.parallel-deserialize-changelog.handler.size | 多线程对变更事件进行解析时,事件处理器的数量。 | 否 | INTEGER | 2 | 说明 仅Flink计算引擎VVR 8.0.11及以上版本支持。 |
metadata-column.include-list | 需要传给下游的元数据列。 | 否 | STRING | 无 | 可用的元数据包括 说明 MySQL CDC YAML连接器无需也不支持添加库名表名和 重要
|
scan.newly-added-table.enabled | 从Checkpoint重启时,是否同步上一次启动时未匹配到的新增表或者移除状态中保存的当前不匹配的表。 | 否 | BOOLEAN | false | 从Checkpoint或Savepoint重启时生效。 重要 在全量阶段时,不支持保存savepoint后,在源表增加新表或删除表再从savepoint重启的操作,这会导致作业无法正常读取数据。 |
scan.binlog.newly-added-table.enabled | 在增量阶段,是否发送匹配到的新增表的数据。 | 否 | BOOLEAN | false | 不能与 |
scan.incremental.snapshot.chunk.key-column | 为某些表指定一列作为快照阶段切分分片的切分列。 | 否 | STRING | 无 |
|
scan.parse.online.schema.changes.enabled | 在增量阶段,是否尝试解析 RDS 无锁变更 DDL 事件。 | 否 | BOOLEAN | false | 参数取值如下:
实验性功能。建议在执行线上无锁变更前,先对Flink作业执行一次快照以便恢复。 说明 仅Flink计算引擎VVR 11.0及以上版本支持。 |
scan.incremental.snapshot.backfill.skip | 是否在快照读取阶段跳过backfill。 | 否 | BOOLEAN | false | 参数取值如下:
backfill仅在单个分片(chunk)快照查询期间生效,不覆盖整个全量读取过程。跳过backfill后,分片快照SQL执行时读到该时刻表的最新数据;分片已读完之后该分片上发生的更新,不再在全量阶段合并,会在进入增量阶段后从Binlog中读取。例如,chunk5快照期间发生的更新会直接体现在chunk5的最新数据中;若已读到chunk80时chunk5才发生更新,该更新会在增量阶段通过Binlog补回。 重要 开启后,分片扫描期间及之后的变更在增量阶段仍会通过Binlog下发,可能与快照数据重复,仅提供at-least-once语义。请确认下游支持按主键幂等写入后再开启。 说明 仅Flink计算引擎VVR 11.1及以上版本支持。 |
treat-tinyint1-as-boolean.enabled | 是否将TINYINT(1)类型当做Boolean类型处理。 | 否 | BOOLEAN | true | 参数取值如下:
|
treat-timestamp-as-datetime-enabled | 是否将TIMESTAMP类型当作DATETIME类型处理。 | 否 | BOOLEAN | false | 参数取值如下:
MySQL TIMESTAMP类型存储的是UTC时间,受时区影响,MySQL DATETIME类型存储的是字面时间,不受时区影响。 开启后会根据server-time-zone将MySQL TIMESTAMP类型数据转换成DATETIME类型。 |
include-comments.enabled | 是否同步表注释和字段注释。 | 否 | BOOELEAN | false | 参数取值如下:
开启后会增加作业内存使用量。 |
scan.incremental.snapshot.unbounded-chunk-first.enabled | 快照读取阶段是否先分发无界的分片。 | 否 | BOOELEAN | false | 参数取值如下:
实验性功能。开启后能够降低TaskManager在快照阶段同步最后一个分片时遇到内存溢出 (OOM) 的风险,建议在作业第一次启动前添加。 说明 仅Flink计算引擎VVR 11.1及以上版本支持。 |
binlog.session.network.timeout | Binlog连接的网络超时时间。 | 否 | DURATION | 10m | 设值为0s时,将会使用MySQL服务端的默认超时时间。 说明 仅Flink计算引擎VVR 11.5及以上版本支持。 |
scan.rate-limit.records-per-second | 限制Source每秒下发的最大记录数。 | 否 | LONG | 无 | 适用于需要限制数据读取场景,此限制在全量和增量阶段都会生效。 Source的 在全量读取阶段,通常需要降低每个批次读取数据的条数进行配合,可以减少 说明 仅Flink计算引擎VVR 11.5及以上版本支持。 |
include-binlog-meta.enable | 是否在消息中携带MySQL Binlog的原始信息,如GTID,Binlog位点等 | 否 | Boolean | false | 适用于原始Binlog同步场景,比如替换原有canal同步链路。 说明 仅Flink计算引擎VVR 11.6及以上版本支持。 |
scan.binlog.tolerate.gtid-holes | 启用此参数可忽略 GTID 序列中的断层,使作业绕过不连续事件并继续运行。 | 否 | Boolean | false | 启用该参数前,必须确保作业的启动位点未过期。若作业从已清理或过期的 GTID 位点启动,引擎将静默跳过缺失的日志,最终导致数据漏读。 说明 仅Flink计算引擎VVR 11.6及以上版本支持此参数。 |
scan.emit.create-table-events.in-batch.enabled | 是否在作业初始化阶段批量下发表结构。 | 否 | Boolean | false | 试验性功能。在单个作业同步表数量较多时,建议开启此选项。 说明 仅Flink计算引擎VVR 11.4及以上版本支持此参数。 |
复用已有 Catalog
自VVR 11.5版本起,您可以在Flink CDC数据摄入作业中直接引用“数据管理”页面中创建的内置MySQL Catalog,减少手写连接属性工作量。
source:
type: mysql
using.built-in-catalog: mysql_rds_catalog目前,数据摄入作业支持自动复用以下 MySQL Catalog 参数:
hostname
port
username
password
catalog.table.metadata-columns
catalog.table.treat-tinyint1-as-boolean
如果希望覆盖以上自动复用的参数,可显式写出相应的 YAML 参数,其具备更高的优先级。
类型映射
数据摄入类型映射如下表所示。
MySQL CDC字段类型 | CDC字段类型 |
TINYINT(n) | TINYINT |
SMALLINT | SMALLINT |
TINYINT UNSIGNED | |
TINYINT UNSIGNED ZEROFILL | |
YEAR | INT |
INT | |
MEDIUMINT | |
MEDIUMINT UNSIGNED | |
MEDIUMINT UNSIGNED ZEROFILL | |
SMALLINT UNSIGNED | |
SMALLINT UNSIGNED ZEROFILL | |
BIGINT | BIGINT |
INT UNSIGNED | |
INT UNSIGNED ZEROFILL | |
BIGINT UNSIGNED | DECIMAL(20, 0) |
BIGINT UNSIGNED ZEROFILL | |
SERIAL | |
FLOAT [UNSIGNED] [ZEROFILL] | FLOAT |
DOUBLE [UNSIGNED] [ZEROFILL] | DOUBLE |
DOUBLE PRECISION [UNSIGNED] [ZEROFILL] | |
REAL [UNSIGNED] [ZEROFILL] | |
NUMERIC(p, s) [UNSIGNED] [ZEROFILL]且p <= 38 | DECIMAL(p, s) |
DECIMAL(p, s) [UNSIGNED] [ZEROFILL]且p <= 38 | |
FIXED(p, s) [UNSIGNED] [ZEROFILL]且p <= 38 | |
BOOLEAN | BOOLEAN |
BIT(1) | |
TINYINT(1) | |
DATE | DATE |
TIME [(p)] | TIME [(p)] |
DATETIME [(p)] | TIMESTAMP [(p)] |
TIMESTAMP [(p)] | 根据
|
CHAR(n) | CHAR(n) |
VARCHAR(n) | VARCHAR(n) |
BIT(n) | BINARY(⌈(n + 7) / 8⌉) |
BINARY(n) | BINARY(n) |
VARBINARY(N) | VARBINARY(N) |
NUMERIC(p, s) [UNSIGNED] [ZEROFILL]且38 < p <= 65 | STRING 说明 在MySQL中,十进制数据类型的精度高达 65,但在Flink中,十进制数据类型的精度仅限于38。所以,如果定义精度大于38的十进制列,则应将其映射到字符串以避免精度损失。 |
DECIMAL(p, s) [UNSIGNED] [ZEROFILL]且38 < p <= 65 | |
FIXED(p, s) [UNSIGNED] [ZEROFILL]且38 < p <= 65 | |
TINYTEXT | STRING |
TEXT | |
MEDIUMTEXT | |
LONGTEXT | |
ENUM | |
JSON | STRING 说明 JSON数据类型将在Flink中转换为JSON格式的字符串。 |
GEOMETRY | STRING 说明 MySQL中的空间数据类型将转换为具有固定JSON格式的字符串,详情请参见MySQL空间数据类型映射。 |
POINT | |
LINESTRING | |
POLYGON | |
MULTIPOINT | |
MULTILINESTRING | |
MULTIPOLYGON | |
GEOMETRYCOLLECTION | |
TINYBLOB | BYTES 说明 对于MySQL中的BLOB数据类型,仅支持长度不大于2147483647(2**31-1)的 blob。 |
BLOB | |
MEDIUMBLOB | |
LONGBLOB |
设置Server ID,避免Binlog消费冲突
数据摄入作业读取Binlog时,以server-id向MySQL注册为复制客户端。若多个作业或其他复制客户端使用相同的Server ID,会导致Binlog消费冲突,作业报错。配置时注意以下事项:
server-id默认随机取5400~6400之间的单个值,多个作业使用默认值时可能冲突,建议显式配置互不重叠的ID。Source并行度大于1时,必须配置Server ID范围,且范围内可用ID数量不小于并行度,每个并行读取器使用不同的ID。
示例:Source并行度为4,配置包含4个ID的范围。
source:
type: mysql
name: MySQL Source
hostname: <hostname>
port: 3306
username: <username>
password: <password>
tables: app_db.\.*
server-id: 5400-5403
sink:
type: hologres多个作业读取同一MySQL实例时,为每个作业分配不重叠的范围。例如作业A使用5400-5403,作业B使用5404-5407。
加速Binlog读取
MySQL连接器作为数据摄入数据源使用时,在增量阶段会解析Binlog文件生成各种变更消息,Binlog文件使用二进制记录着所有表的变更,可以通过以下方式加速Binlog文件解析。
开启并行解析和解析过滤配置
开启配置项
scan.only.deserialize.captured.tables.changelog.enabled:仅对指定表的变更事件进行解析。开启配置项
scan.parallel-deserialize-changelog.enabled:采用多线程对Binlog文件进行解析,并按顺序投放到消费队列。开启该配置时通常需要增加Task Manager CPU进行配合。
优化Debezium参数
debezium.max.queue.size: 162580 debezium.max.batch.size: 40960 debezium.poll.interval.ms: 50debezium.max.queue.size:阻塞队列可以容纳的记录的最大数量。当Debezium从数据库读取事件流时,它会在将事件写入下游之前将它们放入阻塞队列。默认值为8192。debezium.max.batch.size:该连接器每次迭代处理的事件条数最大值。默认值为2048。debezium.poll.interval.ms:连接器应该在请求新的变更事件前等待多少毫秒。默认值为1000毫秒,即1秒。
使用示例:
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: ${mysql.source.table}
server-id: 7601-7604
# Debezium配置
debezium.max.queue.size: 162580
debezium.max.batch.size: 40960
debezium.poll.interval.ms: 50
# 开启解析过滤
scan.only.deserialize.captured.tables.changelog.enabled: trueMySQL CDC 企业版本binlog消费能力为85MB/s,约为开源社区的2倍,当Binlog文件产生速度大于 85MB/s 时(即每6s一个512MB大小的文件),Flink 作业的延迟会持续上升,在Binlog文件产生速度降低后处理延迟会逐步下降。在Binlog文件包含大事务时,可能会导致处理延迟短暂上升,读取完该事务的日志后处理延迟会下降。
分析数据延迟,优化作业吞吐
在增量阶段出现数据延迟时,可以按照以下步骤进行分析:
参见概览中的currentFetchEventTimeLag和currentEmitEventTimeLag两个指标,currentFetchEventTimeLag代表从Binlog读取到数据的延迟,currentEmitEventTimeLag代表从Binlog读取到作业相关的表的数据的延迟。
场景
详情
currentFetchEventTimeLag延迟较小而currentEmitEventTimeLag延迟较大,并且currentEmitEventTimeLag几乎不更新。
currentFetchEventTimeLag延迟较小说明从数据库拉取Binlog的延迟较低,但是Binlog中属于作业需要读取的表的数据较少,因此currentEmitEventTimeLag几乎不更新,属于正常现象。
currentFetchEventTimeLag延迟和currentEmitEventTimeLag延迟都比较大。
说明Source表拉取能力较弱,可以参见本小节的后续步骤进行调优。
反压的存在会导致Source端数据发送至下游算子的速率下降,您可能会观察到sourceIdleTime周期性上升,currentFetchEventTimeLag和currentEmitEventTimeLag不断增长。可以通过增大反压源头所在节点的并发度来避免该情况。
参见CPU中的TM CPU Usage指标和JVM中的TM GC Time指标,确认是否出现CPU或者内存资源不足的情况,可以适当增加作业资源以优化读取性能。
读取RDS归档的OSS日志,避免Binlog过期
使用阿里云RDS MySQL实例作为Source数据源时,支持读取保存在OSS的日志备份。当指定的时间戳或者Binlog位点对应的文件保存在OSS时,会自动拉取OSS日志文件到Flink集群本地进行读取,当指定的时间戳或者Binlog位点对应的文件保存在数据库本地时,会自动切换到使用数据库连接进行读取。该功能仅在实时计算Flink版本提供,社区版MySQL CDC连接器不支持。
开启读取OSS日志备份功能需要配置RDS的连接参数,使用示例:
source:
type: mysql
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: <yourPassword>
tables: <yourTables>
# RDS连接参数,开启读取OSS日志备份
rds.region-id: cn-beijing
rds.access-key-id: your_access_key_id
rds.access-key-secret: your_access_key_secret
rds.db-instance-id: rm-xxxxxxxx # 数据库实例id。
rds.main-db-id: 12345678 # 主库编号。
rds.endpoint: rds.aliyuncs.com使用数据摄入进行整库同步,表结构变更同步
对于只包含数据同步逻辑的作业,建议使用数据摄入运行,数据摄入作业基于数据集成场景进行了深度优化,使用方式参见Flink CDC数据摄入作业以及Flink CDC数据摄入作业开发。
如下代码提供了将MySQL的app_db整库同步到Hologres的示例,对于上游app_db库中的表结构变更,数据摄入作业会将该变更同步到下游数据库:
source:
type: mysql
hostname: <hostname>
port: 3306
username: ${secret_values.mysqlusername}
password: ${secret_values.mysqlpassword}
tables: app_db.\.*
server-id: 5400-5404
sink:
type: hologres
name: Hologres Sink
endpoint: <endpoint>
dbname: <database-name>
username: ${secret_values.holousername}
password: ${secret_values.holopassword}
pipeline:
name: Sync MySQL Database to Hologres数据摄入连接器新增表功能
MySQL的数据摄入连接器针对两种场景下的新增表,分别提供了配置项进行支持。
配置项 | 说明 | 备注 |
| 从Checkpoint重启时,是否同步上一次启动时未匹配到的新增表,全增量同步新增表的数据。 | 仅支持在 |
| 在增量阶段,是否同步匹配到的新增表的数据,自动同步新增表数据。 |
|
在全量阶段时,不支持保存savepoint后,在源表增加新表或删除表再从savepoint重启的操作,这会导致作业无法正常读取数据。
scan.newly-added-table.enabled和scan.binlog.newly-added-table.enabled不建议同时开启,同时开启会导致数据重复问题。