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 |
| 結果テーブルでのデータの更新と削除 | サポートされています |
前提条件
開始する前に、以下を確認してください:
-
AnalyticDB for MySQL クラスターとテーブル。詳細については、「クラスターの作成」および「CREATE TABLE」をご参照ください。
-
クラスターに設定された IP アドレスホワイトリスト。詳細については、「IP アドレスホワイトリストの設定」をご参照ください。
構文
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;