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

DataWorks:MySQL から Elasticsearch へのリアルタイム同期

最終更新日:Jun 23, 2026

データ統合は、MySQL などのソースから Elasticsearch へデータベース全体をリアルタイムで同期することをサポートしています。このトピックでは、MySQL から Elasticsearch へのシナリオを例に、フル同期と増分同期を組み合わせたリアルタイム同期の実行方法について説明します。

前提条件

  • データソースの準備

    • MySQL データソースと Elasticsearch データソースが作成済みであること。詳細については、「データソースの設定」をご参照ください。

    • MySQL データソースで Binlog が有効になっていること。詳細については、「前提条件」をご参照ください。

  • リソースグループ:サーバーレスリソースグループを購入済みであること。

  • ネットワーク接続:リソースグループとデータソース間のネットワーク接続が確立されていること。詳細については、「ネットワーク接続ソリューションの概要」をご参照ください。

タスクの設定

ステップ 1:同期タスクの作成

  1. DataWorks コンソールにログインします。対象のリージョンで、左側のナビゲーションウィンドウにあるデータ統合 > データ統合をクリックします。ドロップダウンリストからワークスペースを選択し、入力 データ統合をクリックします。

  2. 左側のナビゲーションウィンドウで Synchronization Task をクリックします。表示されたページで Create Synchronization Task をクリックし、タスク情報を設定します。

    • Source Type:MySQL。

    • Destination Type:Elasticsearch。

    • Specific Type:データベース全体のリアルタイム同期。

    • Synchronization Mode:

      • Schema Migration:送信先で一致するインデックス構造 (インデックスやフィールドマッピングなど) を自動的に作成します。このステップではデータは移行されません。

      • Full Synchronization (任意):テーブルなどの指定されたソースオブジェクトからすべての既存データを一度に送信先にコピーします。これは通常、初期データ移行や初期化に使用されます。

      • Incremental Sync (オプション): 完全同期が完了した後、ソースからデータの変更 (挿入、更新、削除) を継続的に取得し、送信先に同期します。

ステップ 2:データソースと計算リソースの設定

  1. Source で MySQL データソースを選択します。Destination で Elasticsearch データソースを選択します。

  2. Running Resources セクションで、同期タスクの Resource Group を選択し、タスクに Resource Group の CU を割り当てます。

    説明

    タスクログに Please confirm whether there are enough resources... のようなメッセージが表示された場合、現在のリソースグループで利用可能な計算ユニット (CU) がタスクの開始または実行に不足していることを示します。Configure Resource Group パネルでタスクに割り当てられた CU の数を増やして、より多くの計算リソースを割り当てることができます。

    推奨されるリソースサイズの詳細については、「データ統合の推奨 CU」をご参照ください。実際のご利用要件に基づいて値を調整してください。

  3. ソースデータソースと送信先データソースの両方が 接続チェック に合格することを確認します。

ステップ 3:同期計画の設定

1. データソースの設定

  • このステップでは、[ソーステーブル] セクションでソースデータソースから同期するテーブルを選択し、image アイコンをクリックして右側の [選択されたテーブル] セクションに移動できます。テーブルが多い場合は、Database Filtering または Table filtering を使用して、正規表現を設定することで同期するテーブルを選択できます。

    image

  • 複数のシャーディングされたテーブル (同じスキーマを持つ) からのデータを単一の送信先テーブルに書き込むには、[正規表現でテーブルを選択] できます。

    image
    ソーステーブルの設定で正規表現を入力します。DataWorks は、一致するすべてのソーステーブルを自動的に識別して収集し、そのデータを式によってマッピングされた送信先テーブルに書き込みます。

    説明

    この方法は、シャーディングされたテーブルのマージ同期シナリオ (シャーディングベースの同期に類似) に適用でき、設定効率を向上させ、多数対 1 の同期ルールを繰り返し追加する必要がなくなります。

2. 送信先インデックスマッピングの設定

アクション

説明

Refresh

システムは、選択したソーステーブルを自動的にリストアップします。ただし、送信先インデックスのプロパティは、リフレッシュして確認した後にのみ適用されます。

  • 同期する複数のテーブルを選択し、[マッピング結果を一括リフレッシュ] をクリックします。

  • 送信先インデックス名:送信先インデックス名は、Customize Mapping Rules for Destination Table Names ルールに基づいて自動的に生成されます。デフォルト名は ${source_database_name}_${table_name} です。この名前のインデックスが送信先に存在しない場合、システムは自動的に作成します。

Custom Mapping Rule for Destination Index Name (任意)

システムはデフォルトのルールを使用してインデックス名を生成します:${source_database_name}_${table_name}。また、Customize Mapping Rules for Destination Table Names 列の Edit ボタンをクリックしてカスタムルールを追加することもできます。

  • ルール名:ルールの名前を定義します。ビジネス目的を反映した、わかりやすい名前を指定することを推奨します。

  • 送信先インデックス名:image ボタンをクリックし、Manual Input と Built-in Variable の値を組み合わせて送信先インデックス名を構築できます。サポートされている変数には、ソースデータソース名、ソースデータベース名、ソーステーブル名が含まれます。

  • 組み込み変数の編集:組み込み変数に文字列変換を適用できます。

この機能は、以下のシナリオをサポートします:

  1. 名前にプレフィックスまたはサフィックスを追加する:定数を設定して、ソーステーブル名にプレフィックスまたはサフィックスを追加します。

    ルールの設定

    結果

    [ルール名] フィールドに pre_table_post と入力します。[送信先インデックス名] は、テキスト pre_、[ソーステーブル名] 組み込み変数、およびテキスト _post の 3 つの部分で構成されます。

    pre_table_post ルールを適用すると、ソーステーブル userinfo は送信先インデックス pre_userinfo_post (作成予定) にマッピングされ、ソーステーブル userinfo1 は送信先インデックス pre_userinfo1_post (作成予定) にマッピングされます。

  2. グローバルな文字列置換を実行する:すべてのソーステーブル名に含まれる文字列 dev_ を prd_ に置き換えます。

    ルールの設定

    結果

    [送信先インデックス名] ドロップダウンリストで、[ソーステーブル名] 組み込み変数を選択し、[組み込み変数の編集] ページに移動します。[ソーステーブル名] タブで、ソース文字列 dev_ を送信先文字列 prd_ に置き換える文字列置換ルールを設定します。ルールは上から下に実行されます。[上に移動]、[下に移動]、[削除] を使用して、ルールの順序を調整したり、削除したりできます。

    結果は、ソーステーブル dev_table1 と dev_table2 が、[置き換え] ルールを使用して、送信先インデックス prd_table1 と prd_table2 (作成予定) にマッピングされることを示しています。dev_ プレフィックスは prd_ に置き換えられます。

  3. 複数のテーブルから単一のテーブルにデータを書き込む:送信先インデックス名を定数値に設定します。

    ルールの設定

    結果

    [送信先インデックス名] フィールドに、my_table のような固定値を入力して、すべてのソーステーブルを同じ送信先テーブルにマッピングします。

    結果:ソースデータベース mysql_test2 のソーステーブル table_01 と table_02 は、両方とも同じ送信先インデックス my_table (作成予定) にマッピングされます。これにより、複数のテーブルのデータを単一のインデックスにマージできます。

フィールドデータ型のマッピングを編集 (任意)

システムは、[ソースタイプ] と [送信先タイプ] の間にデフォルトのマッピングを提供します。テーブルの右上隅にある Edit Mapping of Field Data Types をクリックして、ソーステーブルと送信先インデックス間のフィールドデータ型のマッピングをカスタマイズできます。設定が完了したら、Apply and Refresh Mapping をクリックします。

フィールドデータ型のマッピングを編集する際は、型変換ルールが有効であることを確認してください。そうしないと、型変換の失敗が発生し、ダーティデータが生成され、タスクが中断される可能性があります。

送信先インデックスを編集 (任意)

カスタムインデックス名マッピングルールに基づいて、システムは自動的に新しい送信先インデックスを作成するか、一致する名前の既存のインデックスを再利用します。

DataWorks は、ソーステーブルの構造に基づいて送信先インデックスの構造を自動的に生成します。ほとんどの場合、手動での介入は必要ありません。

送信先インデックスのステータスが [作成予定] の場合、元のテーブル構造に基づいて送信先インデックスに新しいフィールドを追加できます。次の操作を実行します:

  1. 送信先インデックスにフィールドを追加します。

    • 単一のインデックスにフィールドを追加する: Destination Index Name 列の image.png アイコンをクリックし、Statement Used to Create Index を編集してフィールドを追加します。

      • Dynamic Mapping Status:データ同期中にソーステーブルから新しいフィールドを送信先インデックスに追加するかどうかを指定します。有効な値:

        • true:システムがソーステーブルに新しいフィールドを検出した場合、そのフィールドを送信先インデックスに追加します。これらのフィールドは検索可能になります。これがデフォルト値です。

        • false:システムがソーステーブルに新しいフィールドを検出した場合、そのフィールドを送信先インデックスに追加しますが、これらのフィールドは検索できません。

        • strict:システムがソーステーブルに新しいフィールドを検出した場合、そのフィールドを送信先インデックスに追加することを拒否し、例外をスローします。エラーの詳細はログで確認できます。

        • runtime:システムがソーステーブルに新しいフィールドを検出した場合、新しいフィールドはインデックスマッピングに追加されません。代わりに、クエリ時にランタイムフィールドとして扱われます。これにより、フィールドをスクリプト計算や検索で使用できます。

        動的マッピングの詳細については、「動的マッピング」をご参照ください。

      • Shards と Replica Shards:インデックスのプライマリシャードとレプリカシャードの数。完全なインデックスは複数のシャードに分割され、異なる Elasticsearch ノードに分散されて、分散検索を可能にし、クエリパフォーマンスを向上させます。詳細については、「基本概念」をご参照ください。

        説明

        Shards と Replica Shards パラメーターの値は、タスク実行後に変更できません。両方のパラメーターのデフォルト値は 1 です。

    • フィールドの一括追加:同期するすべてのテーブルを選択し、テーブルの下部で Batch Modify > Destination Index Structure - Batch Add Fields を選択します。

Value assignment

ネイティブフィールドは、ソースと送信先の間で一致するフィールド名に基づいて自動的にマッピングされます。追加した 新しいフィールドと送信先インデックスのプロパティ には手動で値を割り当てる必要があります。次の操作を実行します:

  • 単一のテーブルに値を割り当てる:Value assignment 列の Configuration ボタンをクリックして、送信先インデックスフィールドに値を割り当てます。

  • 一括で値を割り当てる:リストの下部で Batch Modify > Value assignment を選択して、複数の送信先インデックスにまたがる同一のフィールドに同時に値を割り当てます。

Value Type を変更することで、定数と変数を割り当てることができます。以下のオプションがサポートされています:

  • 送信先インデックスフィールド:

    • 手動割り当て:abc のような定数値を入力します。

    • ソースフィールド:ソーステーブルのフィールドから値を割り当てます。フィールド値または時間値のいずれかを選択できます。

      • フィールド値:ソースフィールドの値を直接送信先に書き込みます。

      • 時間値:ソースフィールドに時間値が含まれている場合、異なるフォーマットを使用して処理し、Destination Format を指定して抽出された値をフォーマットできます。

        • 時間文字列:"2018-10-23 02:13:56" や "2021/05/18" のような時間や日付を表す文字列。時間フォーマットを指定することで、文字列は日付または時間値に解析されます。例えば、上記の例の文字列は yyyy-MM-dd HH:mm:ss と yyyy/MM/dd フォーマットで認識できます。

        • 時間オブジェクト:ソース値が Date や Datetime のような時間データ型の場合、この型を直接選択できます。

        • Unix タイムスタンプ (秒):秒単位の 10 桁のタイムスタンプで、数値または文字列として提供できます。例:1610529203 および "1610529203"。

        • Unix タイムスタンプ (ミリ秒):ミリ秒単位の 13 桁のタイムスタンプで、数値または文字列として提供できます。例:1610529203002 および "1610529203002"。

    • 変数を選択:システムが提供する変数を選択して値のソースとします。

    • 関数:関数を使用してソースフィールドに簡単な変換を適用してから、値として割り当てます。詳細については、「関数式を使用して送信先テーブルフィールドに値を割り当てる」をご参照ください。

  • 送信先インデックスプロパティの割り当て:送信先インデックスのプライマリキーに値を割り当てます。複数のソースフィールドを連結して複合プライマリキーを作成できます。結果の値が一意であることを確認してください。

[ソース分割列]

[ソース分割列] ドロップダウンリストからソーステーブルのフィールドを選択するか、Not Split を選択できます。同期タスクが実行されると、この列に基づいて複数のサブタスクに分割され、データをバッチで並行して読み取ります。

テーブルのプライマリキーをソース分割列として使用することを推奨します。文字列、浮動小数点、日付型はサポートされていません。

ソース分割列は、ソースが MySQL の場合にのみサポートされます。

[フル同期をスキップ]

ステップ 3 でフル同期を設定した場合、個々のテーブルに対してフル同期をスキップすることを選択できます。これは、他の方法で既に完全なデータを送信先に同期している場合に便利です。

Full condition

フル同期フェーズ中にソースデータにフィルターを適用します。WHERE 句のフィルター条件のみを入力してください。WHERE キーワードは含めないでください。

Configure DML Rule

DML メッセージ処理は、ソースからキャプチャされた変更データ (Insert、Update、Delete) が送信先に書き込まれる前に、詳細なフィルタリングと制御を適用します。このルールは増分同期フェーズ中にのみ適用されます。

ステップ 4:詳細設定

詳細パラメーターの設定

タスクをカスタマイズするには、Advanced Parameters タブでパラメーターを変更します。

  1. 右上隅の [詳細設定] をクリックして、詳細パラメーター設定ページに移動します。

  2. 提供されている説明に基づいてパラメーター値を変更します。

  3. AI を活用した設定も利用できます。タスクの同時実行数を調整するコマンドなど、自然言語でコマンドを入力すると、AI モデルが推奨パラメーター値を生成します。AI が生成したパラメーターを受け入れるかどうかを選択できます。

    提案を受け入れるか拒否するかに加えて、[再生成] をクリックして Copilot に新しいパラメーターの推奨を提供させることもできます。

重要

これらのパラメーターは、その目的を完全に理解している場合にのみ変更してください。これにより、タスクの遅延、他のタスクをブロックする過剰なリソース消費、データ損失などの予期しない問題を回避できます。

DDL 機能の設定

一部のリアルタイム同期チャネルは、ソーステーブルスキーマのメタデータ変更を検出し、送信先に更新を同期するように通知したり、アラート、無視、タスクの終了などの他のアクションを実行したりできます。

右上隅の Configure DDL Capability をクリックして、各変更タイプの処理ポリシーを設定できます。サポートされている処理ポリシーはチャネルによって異なります。

  • 通常処理:送信先はソースからの DDL 変更情報を処理します。

  • 無視:変更メッセージは無視され、送信先は変更されません。

  • エラー:リアルタイムのデータベース全体の同期タスクが終了し、ステータスが [エラー] に設定されます。

  • アラート:ソースでこのタイプの変更が発生したときにアラートが送信されます。Configure Alert Rule で DDL 通知ルールを設定する必要があります。

説明

ソースで新しい列が追加され、DDL 同期を通じて送信先で作成された後、システムは送信先テーブルの既存データに対してデータのバックフィルを行いません。

ステップ 5:タスクのデプロイと実行

  1. すべての設定が完了したら、ページ下部の Save をクリックしてタスク設定を保存します。

  2. データベース全体の同期タスクは直接のデバッグをサポートしていません。実行のためには Operation Center にデプロイする必要があります。したがって、新規または編集されたタスクを有効にするには Deploy 操作を実行する必要があります。

  3. デプロイ中に Start immediately after deployment を選択すると、タスクはデプロイと同時に開始されます。そうでない場合、デプロイ後、Data Integration > Synchronization Task に移動し、対象タスクの [操作] 列で手動でタスクを開始する必要があります。

  4. Tasks の対応するタスクの Name/ID をクリックして、タスクの詳細な実行プロセスを表示します。

ステップ 6:アラート設定

1. アラートの作成

Data Integration > Synchronization Task リストで、リアルタイムのデータベース全体のタスクを見つけ、[操作] 列の More > アラーム設定 をクリックして、タスクのアラートポリシーを設定します。

image

(1) Create Rule をクリックしてアラートルールを設定します。

Alert Reason を設定して、Business delay、[フェイルオーバー]、Task status、DDL Notification、Task Resource Utilization などのタスクメトリックを監視し、指定されたしきい値に基づいて CRITICAL または WARNING のアラートレベルを設定できます。

  • Configure Advanced Parameters を設定することで、アラートメッセージ間の時間間隔を制御し、一度に多くのメッセージを送信することによる無駄やメッセージの蓄積を防ぐことができます。

  • アラート理由が Business delay、Task status、または Task Resource Utilization に設定されている場合、タスクが正常に戻ったときに受信者に通知する復旧通知を有効にすることもできます。

(2) アラートルールの管理

既存のアラートルールについては、アラートスイッチを使用してアラートルールを有効または無効にできます。また、アラートレベルに基づいて異なる担当者にアラートを送信することもできます。

2. アラートの表示

タスクリストで More > Configure Alert Rule をクリックしてパネルを展開し、アラートイベントページに移動して、発生したアラートを表示できます。

タスクの管理

タスクの編集

  1. Data Integration > Synchronization Task ページで、作成した同期タスクを見つけます。Operation 列で、More > Edit を選択してタスク情報を変更します。手順は新しいタスクを設定する場合と同じです。

  2. 実行中でないタスクについては、設定を直接変更して保存し、タスクをオペレーションセンターにデプロイして変更を適用できます。

  3. Running のタスクについては、Start immediately after deployment を選択せずにタスクを編集してデプロイすると、元のアクションボタンが Apply Updates に変わります。変更をオペレーションセンターで有効にするには、このボタンをクリックする必要があります。

  4. [更新を適用] をクリックすると、システムはタスクを停止、デプロイ、再起動して変更を適用します。

    • 新しいテーブルを追加または既存のテーブルを切り替える場合:

      更新を適用する際に位置を選択することはできません。更新を確認すると、システムは新しいテーブルに対して スキーマ移行 と フル同期 を実行します。初期化が完了すると、これらのテーブルは元のテーブルと共に増分同期を開始します。

    • その他の情報を変更する場合:

      更新を適用する際に位置を選択できます。確認後、タスクは指定された位置から再開します。位置を指定しない場合、最後に停止した位置 (最後のチェックポイント) から再開します。

    変更されていないテーブルは影響を受けません。更新と再起動後、最後のチェックポイントから再開します。

タスクの表示

同期タスクを作成した後、[同期タスク] ページで作成されたタスクのリストとその基本情報を表示できます。

  • [操作] 列で、同期タスクを Start または Stop できます。[その他] メニューでは、Edit や View などの他の操作を実行できます。

  • 実行中のタスクについては、Execution Overview セクションでそのステータスを表示できます。また、概要の特定のエリアをクリックして実行の詳細を表示することもできます。[表示] をクリックして同期タスクの詳細ページに移動します。上部の [基本情報] セクションには、タスク ID、データソース (例:MySQL_Source → Elasticsearch_Source)、作成時間、同期リソースグループ、ステータス (実行中)、同期計画 (データベース全体のリアルタイム同期)、およびタスクが占有する CU が表示されます。中央の [実行ステータス] セクションでは、プログレスバーを使用して、スキーマ移行、フル同期、リアルタイムデータ同期の 3 つのステージの完了率と実行ステータスが表示されます。

    MySQL から Elasticsearch へのリアルタイム同期タスクは、3 つのステージで構成されます:

    • スキーマ移行:送信先インデックスがどのように作成されたか (既存のインデックスからか、自動作成か) を示します。インデックスが自動作成された場合、DDL ステートメントが表示されます。

    • フル同期:オフライン同期で同期されたテーブル、その進捗状況、および書き込まれたレコード数を表示します。

    • リアルタイムデータ同期:進捗状況、DDL および DML レコード、アラート情報などのリアルタイム統計を表示します。

タスクの再実行

テーブルの追加や削除、送信先テーブルのスキーマやテーブル名情報の変更など、特定のシナリオでは、同期タスクの Operations 列にある Rerun をクリックできます。システムは、新しく追加または変更されたテーブルのみを同期します。以前に同期された、または変更されていないテーブルは再度同期されません。

  • Rerun をクリックして、完全な初期化とリアルタイム同期を再実行します。

  • タスクを編集してテーブルを追加または削除し、タスクを保存してからデプロイします。デプロイ後、[操作] 列に Apply Updates ボタンが表示されます。Apply Updates をクリックして、変更されたタスクの再実行をトリガーします。新しく追加または変更されたテーブルのみが同期されます。以前に同期されたテーブルは再度同期されません。

チェックポイントから再開

ユースケース

タスクの開始位置をリセットすることは、以下のシナリオで役立ちます:

  • タスクの回復とデータの再開:タスクが中断された場合、中断時間を新しい開始位置として手動で指定し、正しいポイントからデータ同期を再開します。

  • データのトラブルシューティングとロールバック:同期後にデータが欠落または異常であることが判明した場合、問題が発生する前の時間に位置をロールバックして、データを再生および修正します。

  • タスク設定の大きな変更:送信先インデックスの構造やフィールドマッピングなど、タスク設定に大幅な調整を行った後、特定の位置から同期を開始するように位置をリセットします。これにより、新しい設定下でのデータの精度が保証されます。

操作手順

Start をクリックします。表示されるダイアログボックスで、Whether to reset the site を選択します。

  • チェックボックスを選択しない場合、タスクは最後に停止したポイント (最後のチェックポイント) から再開します。

  • チェックボックスを選択して時間を指定した場合、タスクは指定された時間から開始します。選択した時間が、ソース Binlog で利用可能な最も古い位置よりも前でないことを確認してください。

重要

無効な位置または存在しない位置に関するエラーが発生した場合は、以下の解決策を使用してください:

  • 位置をリセットする:リアルタイム同期タスクを開始する際に、位置をリセットし、ソースデータベースで利用可能な最も古い位置を選択します。

  • ログの保持期間を調整する:データベースの位置が期限切れになっている場合は、データベースのログ保持期間を、例えば 7 日間に延長します。

  • データを再同期する:データが失われた場合は、再度フル同期を実行するか、オフライン同期タスクを設定して欠落したデータを手動で同期します。

よくある質問

リアルタイムデータベース同期に関するよくある質問については、「データ統合に関するよくある質問」および「データ統合のエラー」をご参照ください。