全部產品
Search
文件中心

Realtime Compute for Apache Flink:MySQL資料攝入連接器

更新時間:Sep 18, 2026

MySQL連接器作為資料來源可以在資料攝入YAML作業中使用。

前提條件

在使用MySQL CDC源表前,必須先按照配置MySQL進行操作,這些操作主要為了滿足使用MySQL CDC源表的前提條件

RDS MySQL

  • 與Realtime ComputeFlink版進行網路探測,確保網路連通。

  • 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

  • 與Realtime ComputeFlink版進行網路探測,確保網路連通。

  • 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

  • 與Realtime ComputeFlink版進行網路探測,確保網路連通。

  • 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

建議填寫Virtual Private Cloud地址。

說明

如果MySQL與即時Flink版不在同一VPC,需要先打通跨VPC的網路或者使用公網的形式訪問,詳情請參見空間管理與操作Flink全託管叢集如何訪問公網?

username

MySQL資料庫服務的使用者名稱。

STRING

無。

password

MySQL資料庫服務的密碼。

STRING

無。

tables

需要同步的MySQL資料表。

STRING

  • 表名支援Regex以讀取多個表的資料。

  • 可以用逗號分隔多個Regex。

說明
  • Regex中請不要使用首尾匹配字元^$。11.2版本中通過點號分割擷取資料庫的Regex,首尾匹配字元會導致擷取到的資料庫Regex不可用。如原來是^db.user_[0-9]+$需要改為db.user_[0-9]+

  • 點號用於分割資料庫名和表名,如果需要用點號匹配任一字元,需要對點號使用反斜線進行轉譯。如:db0.\.*, db1.user_table_[0-9]+, db[1-2].[app|web]order_\.*。

tables.exclude

需要在同步的表中排除的表。

STRING

  • 表名支援Regex以排除多個表的資料。

  • 可以用逗號分隔多個Regex。

說明

點號用於分割資料庫名和表名,如果需要用點號匹配任一字元,需要對點號使用反斜線進行轉譯。如: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

  • 通過英文冒號:串連表名和欄位名,表示一個指定規則,表名可以使用Regex。支援定義多個指定規則,不同指定規則通過英文分號;分割。例如: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位點對應的檔案儲存在資料庫本地時,會自動切換到使用資料庫連接進行讀取。該功能僅在Realtime ComputeFlink版本提供,社區版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

使用資料攝入進行整庫同步,表結構變更同步

對於只包含資料同步邏輯的作業,建議使用資料攝入運行,資料攝入作業基於Data Integration情境進行了深度最佳化,使用方式參見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不建議同時開啟,同時開啟會導致資料重複問題。