全部产品
Search
文档中心

实时计算Flink版:MySQL数据摄入连接器

更新时间:Sep 18, 2026

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(会牺牲复制性能)。

注意事项

数据摄入

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

  • 表名支持正则表达式以读取多个表的数据。

  • 可以用逗号分隔多个正则表达式。

说明
  • 正则表达式中请不要使用首尾匹配字符^$。11.2版本中通过点号分割获取数据库的正则表达式,首尾匹配字符会导致获取到的数据库正则表达式不可用。如原来是^db.user_[0-9]+$需要改为db.user_[0-9]+

  • 点号用于分割数据库名和表名,如果需要用点号匹配任意字符,需要对点号使用反斜杠进行转译。如:db0.\.*, db1.user_table_[0-9]+, db[1-2].[app|web]order_\.*。

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

参数取值如下:

  • initial(默认):在首次启动或无状态启动时,会先扫描历史全量数据,然后读取最新的Binlog数据。

  • latest-offset:在首次启动或无状态启动时,不会扫描历史全量数据,直接从Binlog的末尾(最新的Binlog处)开始读取,即只读取该连接器启动以后的最新变更。

  • earliest-offset:不扫描历史全量数据,直接从可读取的最早Binlog开始读取。

  • specific-offset:不扫描历史全量数据,从您指定的Binlog位点启动,位点可通过同时配置scan.startup.specific-offset.filescan.startup.specific-offset.pos参数来指定从特定Binlog文件名和偏移量启动,也可以只配置scan.startup.specific-offset.gtid-set来指定从某个GTID集合启动。

  • timestamp:不扫描历史全量数据,从指定的时间戳开始读取Binlog。时间戳通过scan.startup.timestamp-millis指定,单位为毫秒。

重要

对于earliest-offsetspecific-offsettimestamp启动模式,如果启动时刻和指定的启动位点时刻的表结构不同,作业会因为表结构不同而报错。换一句话说,使用这三种启动模式,需要保证在指定的Binlog消费位置到作业启动的时间之间,对应表不能发生表结构变更。

scan.startup.specific-offset.file

使用指定位点模式启动时,启动位点的Binlog文件名。

STRING

使用该配置时,scan.startup.mode必须配置为specific-offset。文件名格式例如mysql-bin.000003

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集合格式例如24DA167-0C0C-11E8-8442-00059A3C7B00:1-19

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

该配置生效,需要设置execution.checkpointing.checkpoints-after-tasks-finish.enabled为true。

scan.only.deserialize.captured.tables.changelog.enabled

在增量阶段,是否仅对指定表的变更事件进行反序列化。

BOOLEAN

  • VVR 8.x版本中默认值为false。

  • VVR 11.1及以上版本默认值为true。

参数取值如下:

  • true:仅对目标表的变更数据进行反序列化,加快Binlog读取速度。

  • false(默认):对所有表的变更数据进行反序列化。

scan.parallel-deserialize-changelog.enabled

在增量阶段,是否使用多线程对变更事件进行解析。

BOOLEAN

false

参数取值如下:

  • true:在变更事件的反序列化阶段采用多线程处理,同时保证Binlog事件顺序不变,从而加快读取速度。

  • false(默认):在事件的反序列化阶段使用单线程处理。

说明

仅Flink计算引擎VVR 8.0.11及以上版本支持。

scan.parallel-deserialize-changelog.handler.size

多线程对变更事件进行解析时,事件处理器的数量。

INTEGER

2

说明

仅Flink计算引擎VVR 8.0.11及以上版本支持。

metadata-column.include-list

需要传给下游的元数据列。

STRING

可用的元数据包括op_tses_tsquery_logfilepos,您可以使用英文逗号分隔多个元数据列。

说明

MySQL CDC YAML连接器无需也不支持添加库名表名和op_type元数据列。您可以直接在Transform表达式中使用__data_event_type__来获取变化数据类型,或在Transform表达式中使用__schema_name____table_name__来获取数据库名和表名。

重要
  • file元数据列代表该数据所在的binlog文件,全量阶段为"", 增量阶段为binlog文件名;pos元数据列代表数据所在的binlog文件中的偏移量,全量阶段为"0", 增量阶段为数据在binlog文件中的偏移量,这两个元数据列从:Flink计算引擎VVR 11.5版本开始支持。

  • es_ts元数据列代表changelog在MySQL上对应事务的开始的时间。仅在使用MySQL版本为8.0.x支持,请勿在使用MySQL低版本时添加该元数据列。

  • op_ts时间戳精度到秒,es_ts时间戳精度到毫秒。

scan.newly-added-table.enabled

从Checkpoint重启时,是否同步上一次启动时未匹配到的新增表或者移除状态中保存的当前不匹配的表。

BOOLEAN

false

从Checkpoint或Savepoint重启时生效。

重要

在全量阶段时,不支持保存savepoint后,在源表增加新表或删除表再从savepoint重启的操作,这会导致作业无法正常读取数据。

scan.binlog.newly-added-table.enabled

在增量阶段,是否发送匹配到的新增表的数据。

BOOLEAN

false

不能与scan.newly-added-table.enabled同时开启。

scan.incremental.snapshot.chunk.key-column

为某些表指定一列作为快照阶段切分分片的切分列。

STRING

  • 通过英文冒号:连接表名和字段名,表示一个指定规则,表名可以使用正则表达式。支持定义多个指定规则,不同指定规则通过英文分号;分割。例如:db1.user_table_[0-9]+:col1;db[1-2].[app|web]_order_\\.*:col2

  • 对于无主键表必填,选择的列必须是非空类型(NOT NULL)。有主键的表为选填,仅支持从主键中选择一列。

scan.parse.online.schema.changes.enabled

在增量阶段,是否尝试解析 RDS 无锁变更 DDL 事件。

BOOLEAN

false

参数取值如下:

  • true:解析 RDS 无锁变更 DDL 事件。

  • false(默认):不解析 RDS 无锁变更 DDL 事件。

实验性功能。建议在执行线上无锁变更前,先对Flink作业执行一次快照以便恢复。

说明

仅Flink计算引擎VVR 11.0及以上版本支持。

scan.incremental.snapshot.backfill.skip

是否在快照读取阶段跳过backfill。

BOOLEAN

false

参数取值如下:

  • true:快照读取阶段跳过backfill。

  • false(默认):快照读取阶段不跳过backfill。

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

参数取值如下:

  • true(默认):将TINYINT(1)类型当作Boolean类型处理。

  • false:不将TINYINT(1)类型当作Boolean类型处理。

treat-timestamp-as-datetime-enabled

是否将TIMESTAMP类型当作DATETIME类型处理。

BOOLEAN

false

参数取值如下:

  • true:将MySQL TIMESTAMP类型当作DATETIME类型处理,映射到CDC TIMESTAMP类型。

  • false(默认):将MySQL TIMESTAMP类型映射到CDC TIMESTAMP_LTZ类型。

MySQL TIMESTAMP类型存储的是UTC时间,受时区影响,MySQL DATETIME类型存储的是字面时间,不受时区影响。

开启后会根据server-time-zone将MySQL TIMESTAMP类型数据转换成DATETIME类型。

include-comments.enabled

是否同步表注释和字段注释。

BOOELEAN

false

参数取值如下:

  • true:同步表注释和字段注释。

  • false(默认):不同步表注释和字段注释。

开启后会增加作业内存使用量。

scan.incremental.snapshot.unbounded-chunk-first.enabled

快照读取阶段是否先分发无界的分片。

BOOELEAN

false

参数取值如下:

  • true:快照读取阶段优先分发无界的分片。

  • 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的numRecordsOutPerSecond指标反映整个数据流每秒钟输出的记录数,可以根据这个指标对此参数进行调整。

在全量读取阶段,通常需要降低每个批次读取数据的条数进行配合,可以减少scan.incremental.snapshot.chunk.size参数值。

说明

仅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)]

根据treat-timestamp-as-datetime-enabled参数值,映射字段不同:

true:TIMESTAMP[(p)]

false:TIMESTAMP_LTZ[(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: 50
    • debezium.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: true

MySQL CDC 企业版本binlog消费能力为85MB/s,约为开源社区的2倍,当Binlog文件产生速度大于 85MB/s 时(即每6s一个512MB大小的文件),Flink 作业的延迟会持续上升,在Binlog文件产生速度降低后处理延迟会逐步下降。在Binlog文件包含大事务时,可能会导致处理延迟短暂上升,读取完该事务的日志后处理延迟会下降。

分析数据延迟,优化作业吞吐

在增量阶段出现数据延迟时,可以按照以下步骤进行分析:

  1. 参见概览中的currentFetchEventTimeLag和currentEmitEventTimeLag两个指标,currentFetchEventTimeLag代表从Binlog读取到数据的延迟,currentEmitEventTimeLag代表从Binlog读取到作业相关的表的数据的延迟。

    场景

    详情

    currentFetchEventTimeLag延迟较小而currentEmitEventTimeLag延迟较大,并且currentEmitEventTimeLag几乎不更新。

    currentFetchEventTimeLag延迟较小说明从数据库拉取Binlog的延迟较低,但是Binlog中属于作业需要读取的表的数据较少,因此currentEmitEventTimeLag几乎不更新,属于正常现象。

    currentFetchEventTimeLag延迟和currentEmitEventTimeLag延迟都比较大。

    说明Source表拉取能力较弱,可以参见本小节的后续步骤进行调优。

  2. 反压的存在会导致Source端数据发送至下游算子的速率下降,您可能会观察到sourceIdleTime周期性上升,currentFetchEventTimeLag和currentEmitEventTimeLag不断增长。可以通过增大反压源头所在节点的并发度来避免该情况。

  3. 参见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的数据摄入连接器针对两种场景下的新增表,分别提供了配置项进行支持。

配置项

说明

备注

scan.newly-added-table.enabled

从Checkpoint重启时,是否同步上一次启动时未匹配到的新增表,全增量同步新增表的数据。

仅支持在scan.startup.mode配置项取值为initial模式下使用,其他启动模式下该配置不生效。

scan.binlog.newly-added-table.enabled

在增量阶段,是否同步匹配到的新增表的数据,自动同步新增表数据。

  • 建议在初次启动作业时开启,同步作业会自动解析Create Table DDL并同步数据到下游。如果在数据库表创建结束后,开启该配置重启作业,会导致数据不全问题。

  • 在initial启动模式下,全量阶段结束前所有的DDL操作都无法同步到下游。在全量阶段创建的表,开启了scan.binlog.newly-added-table.enabled也无法完成自动同步。

重要
  • 在全量阶段时,不支持保存savepoint后,在源表增加新表或删除表再从savepoint重启的操作,这会导致作业无法正常读取数据。

  • scan.newly-added-table.enabledscan.binlog.newly-added-table.enabled不建议同时开启,同时开启会导致数据重复问题。