Data Integration は、Kafka や LogHub などの単一テーブルソースから OSS へのリアルタイム同期をサポートしています。本例では、ソースとして Kafka を、送信先として OSS を使用します。
制限事項
Kafka のバージョンは 0.10.2 ~ 2.2.0(両端を含む)の範囲内である必要があります。
前提条件
-
サーバーレスリソースグループまたはデータ統合専用リソースグループを購入済みである必要があります。
-
Kafka データソースおよび OSS データソースを作成済みである必要があります。詳細については、「データソース構成」をご参照ください。
-
リソースグループとデータソース間のネットワーク接続が確立済みである必要があります。詳細については、「ネットワーク接続ソリューションの概要」をご参照ください。
操作手順
ステップ 1:同期タスクタイプの選択
DataWorks コンソールにログインします。対象のリージョンで、左側のナビゲーションウィンドウから をクリックします。ドロップダウンリストからワークスペースを選択し、移動 Data Integration をクリックします。
-
左側のナビゲーションウィンドウで 同期タスク をクリックします。次に、ページ上部の 同期タスクの作成 をクリックします。タスク作成ページで、以下のパラメーターを設定します。
-
Source and destination:
Kafka→OSS -
新規タスク名: 同期タスクのカスタム名を入力します。
-
同期タイプ:
single-table real-time。
-
ステップ 2:ネットワークとリソースの構成
-
ネットワークとリソース構成 セクションで、同期タスクの リソースグループ を選択します。タスクに対して CUs 単位で タスクのリソース使用量 を割り当てることができます。
-
ソースデータソース で、ご利用の
kafkaデータソースを選択します。宛先データソース で、ご利用のOSSデータソースを選択します。その後、接続性のテスト をクリックします。 -
ソースおよび送信先のデータソースの接続性テストが成功したら、次へ をクリックします。
ステップ 3:同期リンクの構成
1. Kafka ソースの構成
ページ上部で Kafka ソースノードをクリックし、Kafka ソース情報 を構成します。

-
Kafka ソース情報 セクションで、Kafka クラスターから同期するトピックを選択します。
その他のパラメーターはデフォルト値を使用するか、必要に応じて変更できます。
-
右上隅で データサンプリング をクリックします。
表示されるダイアログボックスで、開始時間 および サンプリング数 を指定し、収集開始 をクリックします。これにより、指定した Kafka トピックからデータをサンプリングしてプレビューできます。サンプリングされたデータは、後続のデータ処理ノードにおけるデータプレビューおよびビジュアル構成の入力として使用されます。
-
出力フィールド設定 セクションで、同期するフィールドを選択します。
2. データ処理ノードの編集
アイコンをクリックしてデータ処理方法を追加します。利用可能なデータ処理方法は次の 5 種類です:データマスキング、文字列置き換え、データフィルタリング、JSON 解析、および フィールド編集および代入。これらの方法は任意の順序で配置できます。タスク実行時には、構成された順序でデータが逐次処理されます。
各処理ノードの構成後、右上隅の データ出力プレビュー をクリックできます。ダイアログボックスで 上流出力の再取得 をクリックすると、現在のノードがサンプルデータをどのように処理するかをシミュレートし、出力を確認できます。
「Preview Data Output」ダイアログボックスには Input Data および Preview Result の 2 つのセクションがあります。Input Data には、Re-obtain Output Of Ancestor Node をクリックして取得したサンプリングされた Kafka データが表示されます。また、+ Manually Add Data をクリックしてテストデータを追加することも可能です。テーブルには _key_、_value_、_partition_、_offset_、およびタイムスタンプなどのフィールドが含まれます。Preview Result で Preview をクリックすると、処理後の出力およびダーティデータの件数を確認できます。プレビュー結果は参考情報であり、実際のタスク実行結果とは異なる場合があります。
データ出力プレビュー機能は、Kafka ソースの データサンプリング に依存しています。プレビューを実行する前に、Kafka ソース構成でデータサンプリングを実施する必要があります。
3. OSS 送信先の構成
ページ上部で OSS 送信先ノードをクリックし、OSS Destination Information を構成します。
-
OSS Destination Information セクションで、OSS 送信先の基本情報を構成します。
-
書き込みフォーマット: サポートされるフォーマットには Hudi、Paimon、および Iceberg が含まれます。
-
メタデータベース自動構築先の選択: アカウントで Data Lake Formation (DLF) を有効化している場合、データがデータレイクに同期されると、システムが DLF に自動的に対応するメタデータベースおよびメタテーブルを作成します。
説明クロスリージョンでのメタデータ作成はサポートされていません。
-
ストレージパスの選択: データレイクにデータを保存する OSS パスを選択します。
-
ターゲットデータベース: データの送信先データベースを選択します。ライブラリの作成 を選択して DLF メタデータベースを作成し、データベース名 を指定することも可能です。
-
ターゲットテーブル: OSS テーブルにデータを書き込む際に、自動テーブル作成 を選択して自動作成するか、既存のテーブルを使用 を選択して既存のテーブルを使用するかを選択します。
-
テーブル名: 書き込む OSS テーブルの名前を入力または選択します。
-
-
(オプション)テーブルスキーマを編集します。
自動テーブル作成 を選択した場合、テーブルスキーマの編集 ボタンをクリックしてダイアログボックス内でターゲットテーブル構造を編集できます。また、先祖ノードの出力列からテーブルスキーマを再生成します をクリックして、上流ノードの出力カラムから自動的にテーブル構造を生成することも可能です。生成された構造では、プライマリキーとして設定するカラムを選択できます。
-
フィールドマッピングを構成します。
システムは 同名マッピング の原則に基づき、上流カラムとターゲットカラム間のマッピングを自動生成します。必要に応じてこれらのマッピングを調整できます。1 つの上流カラムを複数のターゲットカラムにマッピングできますが、複数の上流カラムを 1 つのターゲットカラムにマッピングすることはできません。上流カラムがターゲットカラムにマッピングされていない場合、そのデータはターゲットテーブルに書き込まれません。
4. アラート
タスクエラーによるビジネスデータ同期の遅延を防ぐため、同期タスクにアラートポリシーを設定できます。
-
ページ右上隅の アラーム設定 をクリックして、サブタスクのリアルタイムアラーム 設定ページを開きます。
-
アラームの作成 をクリックしてアラートルールを構成します。
説明ここで定義したアラートルールは、このタスクによって生成されるリアルタイム同期サブタスクに適用されます。タスク構成後は、「リアルタイム同期タスクの実行と管理」ページで、これらのサブタスクのアラートルールを表示および変更できます。
-
アラートルールを管理します。
既存のアラートルールについて、トグルスイッチを使用して有効化または無効化できます。また、アラートレベルに応じて異なる受信者にアラートを送信できます。
5. 高度な設定
同期タスクには、必要に応じて変更可能な複数のパラメーターが用意されています。
変更を行う前に、各パラメーターの関数を十分に理解し、予期しないエラーやデータ品質の問題を回避してください。
-
ページ右上隅の advanced settings をクリックして高度な設定ページを開きます。
-
advanced settings ページで、必要に応じてパラメーター値を変更します。
ステップ 6:DDL 機能の構成
ソースデータソースにはさまざまな DDL 操作が含まれる可能性があります。ビジネス要件に基づき、ページ右上隅の DDL 能力の設定 をクリックして DDL 機能構成ページに移動し、送信先に同期するさまざまな DDL メッセージに対する処理ポリシーを定義できます。
さまざまな DDL メッセージ処理ポリシーの詳細については、「DDL メッセージ処理ルール」をご参照ください。
ステップ 7:リソースグループの構成
ページ右上隅の リソースグループの設定 をクリックして、タスクの現在のリソースグループを表示および切り替えることができます。
ステップ 8:シミュレーション実行
タスク構成後、ページ右上隅の 模擬実行 をクリックします。この機能は少量のデータサンプルに対してタスク全体をシミュレートし、ターゲットテーブルでの結果をプレビューできます。構成エラー、ランタイム例外、またはダーティデータが存在する場合、リアルタイムでエラーメッセージが表示されます。これにより、タスクが正しく構成され、期待通りの結果を生成することを迅速に検証できます。
-
ダイアログボックスでサンプリングパラメーター(開始時間 および サンプリング数)を設定します。
-
収集開始 をクリックしてサンプルデータを収集します。
-
プレビュー をクリックして、サンプリングされたデータを使用してタスク処理全体をシミュレートします。
ステップ 9:同期タスクの実行
-
すべての設定を完了したら、ページ下部の 設定の完了 をクリックします。
-
ページで、作成したタスクを見つけ、操作 列の 起動 をクリックします。
-
タスク一覧 の対応するタスクの 名前 / ID をクリックして、詳細な実行プロセスを表示します。
同期タスクの管理
タスクステータスの表示
同期タスクを作成後、同期タスクページでタスクステータスおよび詳細を表示できます。
-
起動 列で、タスクを 停止 または 停止 できます。「More」オプションの下には、編集 や 詳細 などの追加操作が用意されています。
-
実行中のタスクについては、実行概要 列で一般的なステータスを確認するか、概要エリアをクリックして詳細な実行情報を表示できます。
Kafka から OSS への単一テーブルリアルタイム同期タスクには、次の 2 つのステージが含まれます。
-
スキーマ移行: 送信先テーブルが既存のテーブルから作成されたか、自動作成されたかを示します。自動作成を選択した場合、システムが使用した DDL 文が表示されます。
-
リアルタイムデータ同期: リアルタイム同期に関する統計情報(ランタイム情報、DDL レコード、アラートなど)を表示します。
タスクの再実行
同期フィールドを変更したり、ターゲットテーブル情報を調整したりする必要がある特殊なケースでは、同期タスクの 操作 列で 再実行 をクリックできます。この操作により、調整されたフィールドおよびその他の変更がターゲットに同期されます。このプロセスでは、変更されていない以前に同期されたテーブルはスキップされます。
-
変更せずに再度タスクを実行するには、再実行 をクリックします。
-
タスクを編集した場合は、変更後に 設定の完了 をクリックします。タスクの操作が アプリケーションの更新 に変更されます。アプリケーションの更新 をクリックすると、新しい構成でタスクが再実行されます。