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

DataWorks:データ同期

最終更新日:Jun 23, 2026

このチュートリアルでは、DataWorks の Data Integration バッチ同期タスクを使用して、異種データソース間でデータをデータウェアハウスに同期する方法について説明します。例として、MySQL の ods_user_info_d テーブル (基本的なユーザー情報) と OSS の user_log.txt ファイル (ウェブサイトのアクセスログデータ) のデータを、StarRocks の ods_user_info_d_starrocks テーブルと ods_raw_log_d_starrocks テーブルに同期します。

前提条件

データを同期する前に、必要な環境が準備されていることを確認してください。詳細については、「環境の準備」をご参照ください。

目的

このチュートリアルの公開データソースから StarRocks にデータを同期します。

ソースタイプ

ソースデータ

ソーススキーマ

送信先タイプ

送信先テーブル

送信先スキーマ

MySQL

テーブル:ods_user_info_d

基本的なユーザー情報

  • uid (ユーザー名)

  • 性別

  • 年齢範囲

  • 十二宮

StarRocks

ods_user_info_d_starrocks

  • uid (ユーザー名)

  • 性別

  • 年齢範囲

  • 十二星座

  • dt (パーティションフィールド)

HttpFile

ファイル:user_log.txt

ユーザーのウェブサイトアクセスログ

各行はユーザーのアクセスレコードを表します。

$remote_addr - $remote_user [$time_local] "$request" $status $body_bytes_sent"$http_referer" "$http_user_agent" [unknown_content];

StarRocks

ods_raw_log_d_starrocks

  • col (生ログ)

  • dt (パーティションフィールド)

DataStudio

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

1. ワークフローの設計

ワークフローの設計

  1. ワークフローの作成

    DataWorks では、すべての開発はワークフロー内で行われます。そのため、ノードを追加する前にワークフローを作成する必要があります。詳細については、「ワークフローの作成」をご参照ください。

    ワークフローに「User Profile Analysis_StarRocks」という名前を付けます。

  2. ワークフローの設計

    ワークフローが作成されると、そのキャンバスが自動的に開きます。「ワークフロー設計」ガイドに従い、キャンバス上の New Node をクリックします。次に、ノードをキャンバスにドラッグし、それらを接続して依存関係を定義することで、データ同期ワークフローを設計します。

    image

  3. このチュートリアルでは、仮想ノードと同期ノードにはデータリネージがないため、キャンバス上でノードを接続してスケジューリング依存関係を定義する必要があります。依存関係を設定する他の方法については、「スケジューリング依存関係の設定ガイド」をご参照ください。次の表に、各ノードのタイプ、名前、および目的を示します。

    ノード分類

    ノードタイプ

    ノード名

    (最終的な出力テーブルにちなんで命名)

    説明

    一般

    仮想ノード

    image

    workshop_start_starrocks

    ユーザープロファイル分析ワークフロー全体を管理し、特に複雑なワークフローにおいてデータフローのパスを明確にします。このノードはドライランタスクであり、コードは不要です。

    データベース

    StarRocks

    image

    ddl_ods_user_info_d_starrocks

    このノードは同期タスクの前に実行してください。MySQL ソースからの基本的なユーザー情報を受け取るための StarRocks テーブル ods_user_info_d_starrocks を作成します。

    データベース

    StarRocks

    image

    ddl_ods_raw_log_d_starrocks

    このノードは同期タスクの前に実行してください。OSS ソースからのユーザーのウェブサイトアクセスレコードを受け取るための StarRocks テーブル ods_raw_log_d_starrocks を作成します。

    Data Integration

    オフライン同期

    image

    ods_user_info_d_starrocks

    MySQL から StarRocks テーブル ods_user_info_d_starrocks に基本的なユーザー情報を同期します。

    Data Integration

    オフライン同期

    image

    ods_raw_log_d_starrocks

    OSS から StarRocks テーブル ods_raw_log_d_starrocks にユーザーのウェブサイトアクセスレコードを同期します。

スケジューリングロジックの設定

このチュートリアルでは、仮想ノード workshop_start_starrocks がワークフロー全体を制御し、毎日 00:30 に実行されるようにスケジュールします。以下に、仮想ノードの主要なスケジューリング設定について説明します。他のノードのスケジューリングを変更する必要はありません。実装ロジックの詳細については、「時間プロパティの詳細設定」をご参照ください。他のスケジューリング設定については、「タスクスケジューリングプロパティ設定の概要」をご参照ください。

スケジューリング設定

設定

説明

スケジュール時刻の設定

[時間プロパティ] セクションで、[スケジュール時刻] を 00:30 に設定します。

仮想ノードは 00:30 に実行されるようにスケジュールされ、ワークフローが毎日実行されるようにトリガーします。

スケジューリング依存関係の設定

[アップストリーム依存関係] セクションで、[ワークスペースのルートノードを使用] チェックボックスをオンにします。

仮想ノード workshop_start_starrocks にはアップストリーム依存関係がないため、Workspace Root Node に直接依存し、それが workshop_start_starrocks ノードの実行をトリガーします。

説明

DataWorks では、すべてのノードにアップストリーム依存関係が必要です。データ同期ステージのすべてのタスクは、仮想ノード workshop_start_starrocks に依存します。したがって、workshop_start_starrocks ノードがデータ同期ワークフローをトリガーします。

2. 同期パイプラインの構築

送信先 StarRocks テーブルの作成

データを同期する前に、受信データを格納するための送信先 StarRocks テーブルを作成します。

この例では、StarRocks テーブルはソーステーブルスキーマに基づいて生成されます。詳細については、「本章の目的」をご参照ください。ワークフローパネルで、データベースノード ddl_ods_user_info_d_starrocks とデータベースノード ddl_ods_raw_log_d_starrocks をダブルクリックし、ノード編集ページに移動して、それぞれに対応する StarRocks テーブル作成コマンドを入力し、image をクリックして保存します。

  • ddl_ods_user_info_d_starrocks

    CREATE TABLE IF NOT EXISTS ods_user_info_d_starrocks (
        uid STRING COMMENT 'ユーザー ID',
        gender STRING COMMENT '性別',
        age_range STRING COMMENT '年齢層',
        zodiac STRING COMMENT '星座',
        dt STRING not null COMMENT '時間'
    ) 
    DUPLICATE KEY(uid)
    COMMENT 'ユーザー行動分析例 - 基本的なユーザー情報テーブル'
    PARTITION BY(dt) 
    PROPERTIES("replication_num" = "1");
  • ddl_ods_raw_info_d_starrocks

    CREATE TABLE IF NOT EXISTS ods_raw_log_d_starrocks (
        col STRING COMMENT 'ログ',
        dt DATE  not null COMMENT '時間'
    ) DUPLICATE KEY(col) 
    COMMENT 'ユーザー行動分析例 - 生のウェブサイトアクセスログテーブル' 
    PARTITION BY(dt) 
    PROPERTIES ("replication_num" = "1");

ユーザーデータのバッチ同期の設定

ワークフローパネルで、バッチ同期ノード ods_user_info_d_starrocks をダブルクリックして ods_user_info_d_starrocks ノードの設定パネルに移動し、例で提供されている MySQL テーブル ods_user_info_d から StarRocks テーブル ods_user_info_d_starrocks に基本的なユーザー情報データを同期するための同期パイプラインを設定します。

  1. ネットワークとリソースの設定

    Data sourceMy Resource Group、および Data going を設定した後、次のステップ をクリックして接続テストを完了します。次の表に設定の詳細を示します。

    パラメーター

    Data source

    • Data source:MySQL

    • Data Source Nameuser_behavior_analysis_mysql

    My Resource Group

    環境準備段階で作成したサーバーレスリソースグループを選択します。

    Data going

    • Data going:StarRocks

    • Data Source NameDoc_StarRocks_Storage_Compute_Tightly_01

  2. タスクの設定

    • ソースと送信先の設定

      モジュール

      パラメーター

      Data source

      Table

      MySQL テーブル ods_user_info_d を選択します。

      Shard Key

      プライマリキーまたはインデックス付きの整数列を分割キーとして使用します。整数型のフィールドのみがサポートされています。

      ここでは、分割キーを uid フィールドに設定します。

      Data going

      Table

      StarRocks テーブル ods_user_info_d_starrocks を選択します。

      Statement Run Before Writing

      この場合、データは dt フィールドによって動的にパーティション分割されます。ノードが再実行されたときに重複したデータ書き込みを防ぐため、各同期の前に次の SQL ステートメントを使用して既存の送信先パーティションを削除します。

      ALTER TABLE ods_user_info_d_starrocks DROP PARTITION IF EXISTS p${var} FORCE、ここで ${var} はパラメーターです。スケジューリングプロパティ設定段階でスケジューリングパラメーターを割り当てることで、スケジューリングシナリオでの動的パラメーター入力を実装できます。詳細については、「スケジューリング設定」をご参照ください。

      Streamload Request Parameters

      StreamLoad のリクエストパラメーターで、JSON 形式である必要があります。

      {
        "row_delimiter": "\\x02",
        "column_separator": "\\x01"
      }
    • フィールドマッピングの設定

      フィールドマッピングは、ソースフィールドと送信先フィールドの関係を定義します。変数にスケジューリングパラメーターを割り当てることで、StarRocks のパーティションフィールドに動的に値を割り当て、日次データが正しいパーティションに書き込まれるようにします。

      [同名フィールドをマッピング] をクリックすると、ソースフィールドが同じ名前の送信先フィールドに自動的にマッピングされます。

      [行を追加] をクリックし、'${var}' を入力し、この値を StarRocks の dt フィールドに手動でマッピングします。

    • スケジューリングプロパティの設定

      設定ページで、右側の [スケジューリングプロパティ] をクリックして Scheduling Configuration パネルを開き、スケジューリングとノード情報を設定します。詳細については、「ノードのスケジューリングプロパティ」をご参照ください。以下のセクションで設定の詳細を説明します。

      パラメーター

      注意

      Scheduling Parameters

      Scheduling Parameters セクションで、Add Parameter をクリックし、以下を追加します:

      • 名前:var

      • 値:$[yyyymmdd-1]

      スケジューリングパラメーターリストには、値が $bizdatebizdate パラメーターも含まれています。

      Scheduling Dependency

      Scheduling Dependency セクションで、このノードの出力としてテーブルを設定します。

      フォーマットは worksspacename.tablename です。

      [このノードの出力名] で、手動で出力名 (例:workspace_name.ods_user_info_d_starrocks) を追加します。この名前は、送信先に設定されたテーブル名と一致する必要があります。

ユーザーログのバッチ同期の設定

ワークフローパネルで、バッチ同期ノード ods_raw_log_d_starrocks をダブルクリックして、ods_raw_log_d_starrocks ノードの設定パネルに入ります。このパネルで、プラットフォームが提供する公開データソースである HttpFile user_log.txt から、StarRocks テーブル ods_raw_log_d_starrocks にユーザーのウェブサイトアクセス情報を転送するための同期リンクを設定します。

  1. ネットワークとリソースの設定

    Data sourceMy Resource Group、および Data going を設定した後、次のステップ をクリックして接続テストを完了します。次の表に設定の詳細を示します。

    パラメーター

    Data source

    • Data source: HttpFile

    • Data Source Nameuser_behavior_analysis_HttpFile

    My Resource Group

    環境準備段階で作成したサーバーレスリソースグループを選択します。

    Data going

    • Data going:StarRocks

    • Data Source NameDoc_StarRocks_Storage_Compute_Tightly_01

  2. タスクの設定

    • ソースと送信先の設定

      モジュール

      パラメーター

      Data source

      File Path

      /user_log.txt

      File Type

      text

      Field Delimiter

      |

      Advanced Settings > Skip Header

      いいえ

      ソースの設定が完了したら、Confirm Data Structure をクリックします。

      Data going

      Table

      ods_raw_log_d_starrocks

      Statement Run Before Writing

      この例では、データは dt フィールドによって動的にパーティション分割されます。ノードが再実行されたときに重複したデータ書き込みを防ぐため、各同期の前に次の SQL ステートメントで既存の送信先パーティションを削除します。

      ALTER TABLE ods_user_info_d_starrocks DROP PARTITION IF EXISTS p${var} FORCE

      ここで、${var} は変数パラメーターです。スケジューリングプロパティを設定する際にスケジューリングパラメーターを割り当てることで、スケジューリングシナリオでの動的パラメーター入力を有効にできます。

      Streamload Request Parameters

      {
        "row_delimiter": "\\x02",
        "column_separator": "\\x01"
      }

      [ソース] モジュールで、user_behavior_analysis_httpfile (HttpFile タイプ) データソースを選択します。詳細設定で、エンコーディングを UTF-8、Null 値を 処理なし、圧縮形式を なし に設定します。[送信先] モジュールで、データソースタイプとして StarRocks を選択します。

    • フィールドマッピングの設定

      ノードツールバーの image アイコンをクリックして、Wizard Mode から Script Mode に切り替えます。これにより、HttpFile ソースのフィールドマッピングを設定し、StarRocks の dt パーティションフィールドに動的に値を割り当てることができます。

      ソース HttpFile の設定で、column 配列に以下を追加します:

       {
                 "type": "STRING",
                "value": "${var}"
              }
    • 以下に、ods_raw_log_d_starrocks ノードの完全なスクリプト例を示します:

      {
          "type": "job",
          "version": "2.0",
          "steps": [
              {
                  "stepType": "httpfile",
                  "parameter": {
                      "fileName": "/user_log.txt",
                      "nullFormat": "",
                      "compress": "",
                      "requestMethod": "GET",
                      "connectTimeoutSeconds": 60,
                      "column": [
                          {
                              "index": 0,
                              "type": "STRING"
                          },
                          {
                              "type": "STRING",
                              "value": "${var}"
                          }
                      ],
                      "skipHeader": "false",
                      "encoding": "UTF-8",
                      "fieldDelimiter": "|",
                      "fieldDelimiterOrigin": "|",
                      "socketTimeoutSeconds": 3600,
                      "envType": 0,
                      "datasource": "user_behavior_analysis",
                      "bufferByteSizeInKB": 1024,
                      "fileFormat": "text"
                  },
                  "name": "Reader",
                  "category": "reader"
              },
              {
                  "stepType": "starrocks",
                  "parameter": {
                      "loadProps": {
                          "row_delimiter": "\\x02",
                          "column_separator": "\\x01"
                      },
                      "envType": 0,
                      "datasource": "Doc_StarRocks_Storage_Compute_Tightly_01",
                      "column": [
                          "col",
                          "dt"
                      ],
                      "tableComment": "",
                      "table": "ods_raw_log_d_starrocks",
                      "preSql": "ALTER TABLE ods_raw_log_d_starrocks DROP PARTITION IF EXISTS  p${var} FORCE ; "
                  },
                  "name": "Writer",
                  "category": "writer"
              },
              {
                  "copies": 1,
                  "parameter": {
                      "nodes": [],
                      "edges": [],
                      "groups": [],
                      "version": "2.0"
                  },
                  "name": "Processor",
                  "category": "processor"
              }
          ],
          "setting": {
              "errorLimit": {
                  "record": "0"
              },
              "locale": "zh",
              "speed": {
                  "throttle": false,
                  "concurrent": 2
              }
          },
          "order": {
              "hops": [
                  {
                      "from": "Reader",
                      "to": "Writer"
                  }
              ]
          }
      }
    • スケジューリングプロパティの設定

      設定ページで、右側の [スケジューリングプロパティ] をクリックして Scheduling Configuration パネルを開き、スケジューリングとノード情報を設定します。以下のセクションで設定について説明します。

      パラメーター

      注意

      Scheduling Parameters

      Scheduling Parameters セクションで、Add Parameter をクリックし、以下を追加します:

      • 名前:var

      • 値:$[yyyymmdd-1]

      Scheduling Dependency

      Scheduling Dependency セクションで、このノードの出力としてテーブルを設定します。

      フォーマットは worksspacename.tablename です。

      Writer ステップのテーブルフィールドの値 (例:ods_raw_log_d_starrocks) が、[このノードの出力名] テーブルに出力名として手動で追加されていることを確認してください。出力名のフォーマットは workspacename.tablename です。

ステップ 3:同期されたデータの検証

ワークフローの実行

  1. ワークフローパネルに移動します。

    ビジネスプロセス で、User Profile Analysis_StarRocks をダブルクリックしてワークフローキャンバスを開きます。image

  2. ワークフローを実行します。

    ワークフローキャンバスで、ツールバーの image アイコンをクリックします。ワークフローのデータ統合ステージのノードが、その依存関係に従って実行されます。

  3. タスクの実行ステータスを確認します。

    image ステータスのノードは、同期が成功したことを示します。

  4. タスク実行ログを表示します。

    キャンバスで、ods_user_info_d_starrocks ノードまたは ods_raw_log_d_starrocks ノードを右クリックし、View Log を選択します。

同期結果の表示

  1. アドホッククエリノードを作成します。

    DataStudio ページの左側のナビゲーションウィンドウで、image をクリックして Ad Hoc Query パネルを開きます。Ad Hoc Query を右クリックし、New Node > スターロックス を選択します。

  2. 送信先テーブルをクエリします。

    -- クエリ内のパーティション列は業務日付に更新する必要があります。例えば、タスクが 20240102 に実行された場合、業務日付はタスク実行日の前日である 20240101 となります。
    SELECT * from ods_raw_log_d_starrocks where dt=your_business_date; 
    SELECT * from ods_user_info_d_starrocks where dt=your_business_date; 

次のステップ

データの同期が完了したので、次のチュートリアルに進み、StarRocks で基本的なユーザー情報とウェブサイトアクセスログを処理する方法を学びます。詳細については、「データの処理」をご参照ください。