AutoETL は、ビジネスシナリオに基づいて同期の動作を調整するための複数の設定オプションを提供しています。たとえば、JSON フィールド変換、検索ルーティング、書き込みモード、リンクで使用されるコンピューティングリソースなどです。このトピックでは、検索ビューと ETL ストアドプロシージャという 2 つの設定エントリポイントのパラメーター構文について説明し、一般的なシナリオでのベストプラクティスを紹介します。
パラメーターは、次の 2 つのエントリポイントで設定できます:
-
検索ビュー: DDL の
WITH (...)句にパラメーターをインラインで設定します。パラメーターは現在の検索ビューにのみ有効です。 -
ETL ストアドプロシージャ: Flink 構文と互換性があります。リンク設定はセッション変数
esl_link_optionsを介して設定でき、同期パラメーター (ソーステーブルの読み取りと宛先への書き込み) は Flink SQL のWITH (...)句に直接書き込まれます。
コンピューティングリソースの設定
ETL リンクは、1 つ以上のワーカースレッドを使用して、実際のデータ転送を実行します。次のパラメーターを使用して、リンクが使用するリソースを制御できます:
|
パラメーター |
説明 |
デフォルト値 |
|
|
現在の ETL リンクの総並列度。 |
4 |
|
|
各ワーカーに割り当てられる CPU コア数。 |
2 |
|
|
各ワーカーがサポートするタスクの同時実行数。 |
4 |
ワーカーの同時実行数は、高ければ高いほど良いというわけではありません。すべてのワーカーの合計 CPU 容量も考慮する必要があります。
転送リンクは CU で測定されます。CU の計算式は CU = parallelism / link.tm.slot * link.tm.cpu + 1 です。次の設定例は 5 CU に相当します。
検索ビュー
CREATE SEARCH VIEW view_test
WITH (
'parallelism' = '8',
'link.tm.cpu' = '4',
'link.tm.slot' = '8'
) AS SELECT * FROM t1;
ETL ストアドプロシージャ
SET esl_link_options = "'parallelism' = '8', 'link.tm.cpu' = '4', 'link.tm.slot' = '8'";
CALL dbms_etl.sync_by_sql("search", "
CREATE TEMPORARY TABLE `db1`.`t1` (
`id` BIGINT,
`c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'database-name' = 'db1',
'table-name' = 't1'
);
CREATE TEMPORARY TABLE `dest` (
`id` BIGINT,
`c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'opensearch',
'index' = 'dest'
);
INSERT INTO `dest`
SELECT * FROM `db1`.`t1`;
");
サポートされるパラメーター一覧
次の表は、現在 AutoETL でサポートされているすべてのパラメーターをまとめたものです。
リンクパラメーター
|
パラメーター |
説明 |
デフォルト値 |
|
|
現在の ETL リンクの総並列度。 |
4 |
|
|
各ワーカーに割り当てられる CPU コア数。 |
2 |
|
|
各ワーカーがサポートするタスクの同時実行数。 |
4 |
|
|
ETL リンクのチェックポイント間隔。 |
180 s |
同期パラメーター
|
パラメーター |
説明 |
デフォルト値 |
|
|
同期先の GDN クラスター。 |
現在のインスタンスのデフォルトクラスター |
|
|
PolarSearch における検索ビューのインデックス名。ストアドプロシージャの場合、宛先テーブルは |
検索ビュー名 |
|
|
フルスキャンフェーズにおける各チャンクの行数。このパラメーターは、フルスライシングの粒度と同時実行効率に影響します。 |
131072 |
|
|
ドキュメントを PolarSearch の指定されたシャードにルーティングするために使用されるフィールド名。 |
なし |
|
|
削除操作を無視するかどうか。 |
false |
|
|
Index モードで強制的に行全体の置換書き込みを行うかどうか。そうでない場合は、フィールドの部分更新に Update モードが使用されます。 |
false |
|
|
JSON パースが必要なフィールド名。 |
なし |
|
|
JSON パースモード: |
nested |
ベストプラクティス
JSON フィールドの自動変換
デフォルトでは、MySQL の JSON フィールドは、PolarSearch に同期される際に文字列として保存されます。PolarDB の検索ビューでは、PolarSearch への同期中に MySQL の JSON フィールドを自動的にパースし、ネストされたフィールドを生成したり、JSON フィールドをフラット化したりできます。
-
データ準備
クラスターで次の SQL ステートメントを実行して、サンプルデータベースとテーブルを作成し、テストデータを挿入します:
CREATE DATABASE IF NOT EXISTS db1; USE db1; CREATE TABLE IF NOT EXISTS t1 ( id INT PRIMARY KEY, c1 JSON ); INSERT INTO t1(id, c1) VALUES (1, '{"age": 75, "name": "User_A5pqo", "tags": ["q7XG", "Unx9", "EBy8"], "active": false, "metadata": {"source": "script", "version": "1.0", "created_at": "2026-03-06T07:04:45.264573Z"}}'), (2, '{"age": 55, "name": "User_xL1YH", "tags": ["QNcC", "kqU7"], "active": true, "metadata": {"source": "script", "version": "1.0", "created_at": "2026-03-06T07:04:45.264632Z"}}'), (3, '{"age": 25, "name": "User_zoRSH", "tags": [ ], "active": true, "metadata": {"source": "script", "version": "1.0", "created_at": "2026-03-06T07:04:45.264654Z"}}'); -
デフォルト設定 (文字列ストレージ)
検索ビュー
CREATE SEARCH VIEW json_test AS SELECT * FROM t1;ETL ストアドプロシージャ
CALL dbms_etl.sync_by_sql("search", " CREATE TEMPORARY TABLE `db1`.`t1` ( `id` INT, `c1` STRING, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'mysql', 'database-name' = 'db1', 'table-name' = 't1' ); CREATE TEMPORARY TABLE `dest` ( `id` INT, `c1` STRING, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'opensearch', 'index' = 'json_test' ); INSERT INTO `dest` SELECT * FROM `db1`.`t1`; ");データ検証
{ "_index" : "json_test", "_id" : "3", "_source" : { "id" : 3, "c1" : "{\"age\":25,\"name\":\"User_zoRSH\",\"tags\":[ ],\"active\":true,\"metadata\":{...}}" } } -
c1 をネスト型に変換
検索ビュー
CREATE SEARCH VIEW json_test WITH ( 'sink.json-flatten.fields' = 'c1' -- ; で複数のフィールドを区切ります ) AS SELECT * FROM t1;ETL ストアドプロシージャ
CALL dbms_etl.sync_by_sql("search", " CREATE TEMPORARY TABLE `db1`.`t1` ( `id` INT, `c1` STRING, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'mysql', 'database-name' = 'db1', 'table-name' = 't1' ); CREATE TEMPORARY TABLE `dest` ( `id` INT, `c1` STRING, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'opensearch', 'index' = 'json_test', 'sink.json-flatten.fields' = 'c1' -- ; で複数のフィールドを区切ります ); INSERT INTO `dest` SELECT * FROM `db1`.`t1`; ");データ検証
{ "_index" : "json_test", "_id" : "3", "_source" : { "id" : 3, "c1" : { "age" : 25, "name" : "User_zoRSH", "tags" : [ ], "active" : true, "metadata" : { "source" : "script", "version" : "1.0", "created_at" : "2026-03-06T07:04:45.264654Z" } } } } -
c1 を JSON の第1レベルでフラット化
検索ビュー
CREATE SEARCH VIEW json_test WITH ( 'sink.json-flatten.fields' = 'c1', 'sink.json-flatten.mode' = 'flatten' ) AS SELECT * FROM t1;ETL ストアドプロシージャ
CALL dbms_etl.sync_by_sql("search", " CREATE TEMPORARY TABLE `db1`.`t1` ( `id` INT, `c1` STRING, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'mysql', 'database-name' = 'db1', 'table-name' = 't1' ); CREATE TEMPORARY TABLE `dest` ( `id` INT, `c1` STRING, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'opensearch', 'index' = 'json_test', 'sink.json-flatten.fields' = 'c1', 'sink.json-flatten.mode' = 'flatten' ); INSERT INTO `dest` SELECT * FROM `db1`.`t1`; ");データ検証: フラット化後、JSON の第1レベルのフィールドは、PolarSearch インデックス内で独立したトップレベルフィールドになります。
{ "_index" : "json_test", "_id" : "3", "_source" : { "id" : 3, "metadata" : { "source" : "script", "version" : "1.0", "created_at" : "2026-03-06T07:04:45.264654Z" }, "name" : "User_zoRSH", "active" : true, "age" : 25, "tags" : [ ] } }
検索ルーティングフィールドの設定
検索ビューで 1 つ以上のフィールド名を指定して、データ行を PolarSearch の指定されたシャードにルーティングします。
-
データ準備
CREATE DATABASE IF NOT EXISTS db1; USE db1; CREATE TABLE IF NOT EXISTS t1 ( id INT PRIMARY KEY, c1 BIGINT ); INSERT INTO t1(id, c1) VALUES (1, 3), (2, 2), (3, 1); -
ルーティングフィールドの設定
検索ビュー
CREATE SEARCH VIEW routing_test WITH ( 'routing-fields' = 'c1' -- ; で複数のフィールドを区切ります ) AS SELECT * FROM t1;ETL ストアドプロシージャ
CALL dbms_etl.sync_by_sql("search", " CREATE TEMPORARY TABLE `db1`.`t1` ( `id` INT, `c1` BIGINT, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'mysql', 'database-name' = 'db1', 'table-name' = 't1' ); CREATE TEMPORARY TABLE `dest` ( `id` INT, `c1` BIGINT, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'opensearch', 'index' = 'routing_test', 'routing-fields' = 'c1' -- ; で複数のフィールドを区切ります ); INSERT INTO `dest` SELECT * FROM `db1`.`t1`; ");データ検証: ドキュメントは
c1フィールドの値に基づいて、対応するシャードにルーティングされます。_routingフィールドには、ルーティングに使用された値が記録されます。{ "_index" : "routing_test", "_id" : "1", "_routing" : "3", "_source" : { "id" : 1, "c1" : 3 } }, { "_index" : "routing_test", "_id" : "3", "_routing" : "1", "_source" : { "id" : 3, "c1" : 1 } }, { "_index" : "routing_test", "_id" : "2", "_routing" : "2", "_source" : { "id" : 2, "c1" : 2 } }
削除の無視
複数のテーブルを集約する検索ビューの場合、AutoETL は削除してから挿入することで宛先インデックスを更新します。削除されたデータの中間状態にクエリがアクセスしないようにしたい場合は、削除の無視を有効にして、同期中にリンクが削除操作をスキップするように設定できます。
検索ビュー
CREATE SEARCH VIEW view_test
WITH (
'ignore-delete' = 'true'
) AS SELECT * FROM t1;
ETL ストアドプロシージャ
CALL dbms_etl.sync_by_sql("search", "
CREATE TEMPORARY TABLE `db1`.`t1` (
`id` BIGINT,
`c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'database-name' = 'db1',
'table-name' = 't1'
);
CREATE TEMPORARY TABLE `dest` (
`id` BIGINT,
`c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'opensearch',
'index' = 'view_test',
'ignore-delete' = 'true'
);
INSERT INTO `dest`
SELECT * FROM `db1`.`t1`;
");
この設定後、リンクは削除操作を実行しなくなります。
データがクリーンアップされないため、PolarSearch インデックスが大きくなる可能性があります。MySQL ソーステーブルのフィールドを使用して、削除された行をマークすることを推奨します。PolarSearch に同期した後、スケジュールされたタスクを使用して、マークされたドキュメントをクリーンアップできます。
置換書き込み
AutoETL が PolarSearch インデックスに書き込む際、ドキュメントがすでに存在する場合、デフォルトで Update モードを使用します。このモードでは、検索ビューによって書き込まれたフィールドのみが更新されます。行全体の置換が必要なシナリオでは、sink.force-index-request を設定して、Index 書き込みモードを有効にできます。
検索ビュー
CREATE SEARCH VIEW view_test
WITH (
'sink.force-index-request' = 'true'
) AS SELECT * FROM t1;
ETL ストアドプロシージャ
CALL dbms_etl.sync_by_sql("search", "
CREATE TEMPORARY TABLE `db1`.`t1` (
`id` BIGINT,
`c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'database-name' = 'db1',
'table-name' = 't1'
);
CREATE TEMPORARY TABLE `dest` (
`id` BIGINT,
`c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'opensearch',
'index' = 'view_test',
'sink.force-index-request' = 'true'
);
INSERT INTO `dest`
SELECT * FROM `db1`.`t1`;
");