すべてのプロダクト
Search
ドキュメントセンター

Realtime Compute for Apache Flink:AnalyticDB for MySQL V3.0

最終更新日:Aug 26, 2026

AnalyticDB for MySQL V3.0 コネクタを使用すると、Flink SQL を使用して AnalyticDB for MySQL V3.0 クラスターの読み取りと書き込みができます。AnalyticDB for MySQL は、高スループットのリアルタイム書き込み、低レイテンシーの分析、および複雑な ETL (抽出・変換・書き出し) 操作をサポートするクラウドネイティブのデータウェアハウスサービスです。

このコネクタは、次のテーブルタイプと機能をサポートしています:

項目 説明
テーブルタイプ ソーステーブル、ディメンションテーブル、結果テーブル。
説明

ソーステーブルには Ververica Runtime (VVR) 8.0.4 以降が必要です。ソーステーブルのパラメーターについては、「Flink を使用したバイナリログのサブスクライブ」をご参照ください。

実行モード ストリーミングモードとバッチモード
データフォーマット N/A
メトリクス N/A
API タイプ SQL API
結果テーブルでのデータの更新と削除 サポートされています

前提条件

開始する前に、以下を確認してください:

構文

CREATE TEMPORARY TABLE adb_table (
  `id` INT,
  `num` BIGINT,
  PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
  'connector' = 'adb3.0',
  'url' = 'jdbc:mysql://<endpoint>:<port>/<databaseName>',
  'userName' = '<yourUsername>',
  'password' = '<yourPassword>',
  'tableName' = '<yourTablename>'
);
重要

Flink DDL 文で定義されたプライマリキーは、AnalyticDB for MySQL の物理テーブルのプライマリキーと一致する必要があります。つまり、フィールド、名前が同じで、両方に存在する必要があります。プライマリキーが一致しないと、データ破損の原因となる可能性があります。

主な動作

プライマリキーと書き込みモード

DDL でプライマリキーが定義されている場合、コネクタはアップサートモードで動作します。このモードでは、重複するプライマリキーは新しい挿入ではなく更新をトリガーします。具体的な SQL の動作は replaceMode パラメーターによって決まります:

  • replace — REPLACE INTO を使用します。プライマリキーが重複する場合、既存の行全体が上書きされます。

  • upsert — INSERT INTO ... ON DUPLICATE KEY UPDATE を使用します。指定されたフィールドのみが更新され、他のフィールドは現在の値を保持します。

  • insert — INSERT IGNORE INTO を使用します。重複するプライマリキーは警告なしに無視され、既存の行が保持されます。

プライマリキーが定義されていない場合、コネクタは常に INSERT IGNORE INTO を使用します。

replaceMode は、文字列値 (replace、upsert、insert) を使用する場合、AnalyticDB for MySQL V3.1.3.5 以降および VVR 11.2 以降が必要です。以前の VVR バージョンでは、true (replace に相当) と false (upsert に相当) を使用します。VVR 11.2 以降は、true および false との互換性を維持しています。

書き込みバッファリング

コネクタはレコードをメモリにバッファリングし、バッチでフラッシュします。フラッシュは、以下のいずれかの条件が満たされたときにトリガーされます:

  • バッファリングされたレコード数が batchSize (デフォルト:1,000) または bufferSize (デフォルト:1,000) に達したとき

  • 最後のフラッシュからの時間が flushIntervalMs (デフォルト:3,000 ms) に達したとき

batchSize と bufferSize は、プライマリキーが定義されている場合にのみ有効です。

ディメンションテーブルのキャッシュ

コネクタは、ディメンションテーブルに対して 3 つのキャッシュポリシーをサポートしています。キャッシュを使用すると、物理テーブルに対するルックアップが減少しますが、追加のメモリが必要になります。

キャッシュポリシー 動作 使用場面
ALL (デフォルト) ジョブ開始前にすべてのディメンションテーブルデータがメモリにロードされます。その後のルックアップはキャッシュのみにヒットします。データは cacheTTLMs の有効期限が切れた後に再読み込みされます。 ディメンションテーブルが小さく、存在しないキーが頻繁にある場合
LRU 頻繁にアクセスされる行がキャッシュされます。キャッシュミスが発生した場合、コネクタは物理テーブルをクエリしてキャッシュを更新します。エントリは cacheTTLMs の後に有効期限切れになります。 ディメンションテーブルは大きいが、アクセスがキーのサブセットに偏っている場合
None キャッシュなし。すべてのルックアップで物理テーブルを直接クエリします。 ルックアップの量が少ない、またはデータの新鮮さに関する要件が厳しい場合
重要

ALL キャッシュを使用する場合、コネクタは非同期でディメンションテーブル全体をメモリにロードします。メモリ不足 (OOM) エラーを回避するために、テーブル結合に使用されるノードのメモリを増やす必要があります。必要な増加量は、リモートテーブルのサイズの少なくとも 2 倍です。ジョブ中はノードのメモリ使用量を監視してください。

WITH 句のパラメーター

共通パラメーター

パラメーター 型 必須 デフォルト 説明
connector String はい — adb3.0
url String はい — データベースの Java Database Connectivity (JDBC) URL。フォーマットは jdbc:mysql://<endpoint>:<port>/<databaseName> です。エンドポイントとポートは、AnalyticDB for MySQL コンソール のクラスターページにある [ネットワーク情報]
userName String はい — データベースアクセス用のユーザー名
password String はい — データベースアクセス用のパスワード
tableName String はい — データベース内のターゲットテーブルの名前
maxRetryTimes Integer いいえ 10 読み取りまたは書き込みの失敗における最大リトライ回数

シンク テーブル パラメーター

パラメーター 型 必須 デフォルト 説明
batchSize Integer いいえ 1000 バッチごとに書き込まれるレコード数。プライマリキーが定義されている場合にのみ有効です。
bufferSize Integer いいえ 1000 フラッシュがトリガーされる前にメモリにバッファリングされるレコードの最大数。プライマリキーが定義されている場合にのみ有効です。
flushIntervalMs Integer いいえ 3000 フラッシュ間の最大時間 (ミリ秒)。この間隔が経過すると、バッファーサイズに関係なく、バッファリングされたすべてのレコードが書き込まれます。
ignoreDelete Boolean いいえ false 削除操作を無視するかどうか。削除をスキップするには true に設定し、適用するには false に設定します。
replaceMode Boolean いいえ true プライマリキーが定義されている場合の書き込みモード。VVR 11.2 以降:replace、upsert、または insert。以前のバージョン:true (replace と同じ) または false (upsert と同じ)。AnalyticDB for MySQL V3.1.3.5 以降が必要です。DDL でプライマリキーが定義されている場合にのみ有効です。それ以外の場合は、常に INSERT IGNORE INTO が使用されます。
excludeUpdateColumns String いいえ (空) replaceMode が upsert (または false) の場合に更新から除外する列のカンマ区切りリスト。プライマリキーが重複する場合:残りの列のみが更新され、除外された列は既存の値を保持します。プライマリキーが一意の場合:すべての列が挿入されます。例:excludeUpdateColumns='column1,column2'。値は改行を含まない単一行である必要があります。
connectionMaxActive Integer いいえ 40 スレッドプール内の同時データベース接続の最大数

ディメンションテーブルのパラメーター

パラメーター 型 必須 デフォルト 説明
cache String いいえ ALL キャッシュポリシー:None、LRU、または ALL。詳細については、「ディメンションテーブルのキャッシュ」をご参照ください。
cacheSize Integer いいえ 100000 キャッシュされる行の最大数。cache が LRU の場合に必須です。
cacheTTLMs Integer いいえ Long.MAX_VALUE キャッシュエントリの有効期間 (ミリ秒)。LRU の場合:この期間後にエントリは有効期限切れになります (デフォルト:有効期限なし)。ALL の場合:この期間後にキャッシュ全体が再読み込みされます (デフォルト:再読み込みなし)。cache が None の場合は使用されません。
maxJoinRows Integer いいえ 1024 入力レコードごとに一致するディメンションテーブルの行の最大数。各入力レコードが最大 n 個のディメンションテーブルの行に一致する場合、これを n に設定して Realtime Compute for Apache Flink での結合パフォーマンスを最適化します。

データ型のマッピング

AnalyticDB for MySQL V3.0 Realtime Compute for Apache Flink
BOOLEAN BOOLEAN
TINYINT TINYINT
SMALLINT SMALLINT
INT INT
BIGINT BIGINT
FLOAT FLOAT
DOUBLE DOUBLE
DECIMAL(p, s) or NUMERIC(p, s) DECIMAL(p, s)
VARCHAR STRING
BINARY BYTES
DATE DATE
TIME TIME
DATETIME TIMESTAMP
TIMESTAMP TIMESTAMP
POINT STRING

例

結果テーブル

次の例では、datagen ソースからデータを読み取り、AnalyticDB for MySQL の結果テーブルに書き込みます。

CREATE TEMPORARY TABLE datagen_source (
  `name` VARCHAR,
  `age`  INT
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE adb_sink (
  `name` VARCHAR,
  `age`  INT
) WITH (
  'connector' = 'adb3.0',
  'url'       = 'jdbc:mysql://<endpoint>:<port>/<databaseName>',
  'userName'  = '<yourUsername>',
  'password'  = '<yourPassword>',
  'tableName' = '<yourTablename>'
);

INSERT INTO adb_sink
SELECT * FROM datagen_source;

ディメンションテーブル

次の例では、datagen ソースと AnalyticDB for MySQL のディメンションテーブルをテンポラル結合で結合します。結果は blackhole の結果テーブルに書き込まれます。

CREATE TEMPORARY TABLE datagen_source (
  `a`       INT,
  `b`       VARCHAR,
  `c`       STRING,
  `proctime` AS PROCTIME()
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE adb_dim (
  `a` INT,
  `b` VARCHAR,
  `c` VARCHAR
) WITH (
  'connector' = 'adb3.0',
  'url'       = 'jdbc:mysql://<endpoint>:<port>/<databaseName>',
  'userName'  = '<yourUsername>',
  'password'  = '<yourPassword>',
  'tableName' = '<yourTablename>'
);

CREATE TEMPORARY TABLE blackhole_sink (
  `a` INT,
  `b` VARCHAR
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT T.a, H.b
FROM datagen_source AS T
JOIN adb_dim FOR SYSTEM_TIME AS OF T.proctime AS H
  ON T.a = H.a;

リファレンス