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

E-MapReduce:Spark Load

最終更新日:Mar 27, 2026

Spark Load は、StarRocks にデータをインポートする前に、外部の Spark クラスターにデータの前処理をオフロードします。これにより、StarRocks クラスターの計算負荷が軽減され、初回移行や TB レベルの大規模なインポートに推奨されるアプローチです。

重要

ユーザーアカウントは、インポートジョブを送信する前に、Spark リソースに対する USAGE-PRIV 権限を持っている必要があります。

仕組み

Spark Load は非同期のインポートメソッドです。MySQL プロトコルを使用してジョブを送信し、SHOW LOAD で結果を確認します。

以下の図は、そのワークフローを示しています。

Spark Load
  1. フロントエンドノードに Spark Load ジョブを送信します。

  2. フロントエンドノードは、抽出、変換、ロード (ETL) ジョブをスケジュールし、Spark クラスターに送信します。

  3. Spark クラスターは ETL ジョブを実行します。ビットマップグローバル辞書を構築し、データのパーティション分割、ソート、集約を行います。

  4. ETL ジョブが完了すると、フロントエンドノードは各パーティションの前処理済みデータディレクトリを特定し、バックエンドノードにプッシュジョブを実行するようスケジュールします。

  5. バックエンドノードは Broker を使用して Hadoop 分散ファイルシステム (HDFS) からデータを読み取り、StarRocks が内部で保存するフォーマットに変換します。

  6. フロントエンドノードは新しい StarRocks バージョンを公開し、インポートジョブを完了としてマークします。

基本概念

Spark ETL:インポート前にデータに対して ETL 操作を実行する Spark プログラムです。グローバル辞書の作成、データのパーティション分割、ソート、集約を処理します。

Broker:ファイルシステムインターフェイスをラップするステートレスなプロセスで、StarRocks が HDFS などのリモートストレージシステムからファイルを読み取れるようにします。

グローバル辞書:生の値とエンコードされた整数をマッピングするデータ構造です。グローバル辞書は、インポート前にビットマップ列を事前計算するために使用されます。StarRocks のビットマップ列は roaring bitmap を使用しており、整数の入力が必要です。

グローバル辞書のワークフロー

StarRocks のビットマップ列には整数の入力が必要です。グローバル辞書がこの変換を処理します。

  1. ソースデータを読み取り、一時的な Hive テーブル (hive-table) に保存します。

  2. hive-table から値を重複排除し、distinct-value-table という名前のテーブルに格納します。

  3. 辞書テーブル (dict-table) を作成します。このテーブルには、生の値用の列とエンコードされた整数用の列が 1 つずつあります。

  4. distinct-value-tabledict-table を LEFT JOIN し、ウィンドウ関数を適用して新しい生の値をエンコードし、結果を dict-table に書き戻します。

  5. dict-tablehive-table を JOIN して、生の値をエンコードされた整数に置き換えます。

  6. エンコードされた hive-table データを後続の ETL ステップに渡し、StarRocks にインポートします。

重要

グローバル辞書は、Hive テーブルからインポートする場合にのみサポートされます。

データの前処理

グローバル辞書が構築された後 (該当する場合)、Spark はデータを前処理します。

  1. HDFS ファイルまたは Hive テーブルからデータを読み取ります。

  2. フィールドマッピングと式ベースの計算を適用します。パーティション情報に基づいて bucket-id フィールドを生成します。

  3. StarRocks テーブルのロールアップメタデータからロールアップツリーを構築します。

  4. ロールアップツリーを走査し、データをレイヤーごとに集約します。各レイヤーは前のレイヤーから計算されます。

  5. bucket-id に基づいて集約データをバケットに分散し、HDFS に書き込みます。

  6. Broker は HDFS ファイルを読み取り、データを StarRocks のバックエンドノードにプッシュします。

前提条件

開始する前に、以下が揃っていることを確認してください。

  • 実行中の Spark クラスターと YARN ResourceManager を備えた EMR クラスター

  • StarRocks で設定された Broker (ALTER SYSTEM ADD BROKER)

  • Spark 外部リソースに付与された USAGE-PRIV 権限

  • Spark 2.4.5 またはそれ以降の 2.x バージョンがダウンロードされ、フロントエンドノードに保存されていること

  • Hadoop 2.5.2 またはそれ以降の 2.x バージョンがダウンロードされ、フロントエンドノードに保存されていること

Spark Load の設定

外部 Spark リソースの作成 → Spark クライアントの設定 → YARN クライアントの設定 → インポートジョブの送信、の順にこれらのステップを完了します。

ステップ 1:外部 Spark リソースの作成

Spark クラスターを StarRocks の外部リソースとして登録し、フロントエンドノードが ETL ジョブを送信できるようにします。

構文

CREATE EXTERNAL RESOURCE "resource_name"
PROPERTIES
(
    "type" = "spark",
    "spark.master" = "yarn",
    "spark.submit.deployMode" = "<cluster|client>",
    "spark.hadoop.fs.defaultFS" = "hdfs://<namenode-host>:<port>",
    "spark.hadoop.yarn.resourcemanager.address" = "<resourcemanager-host>:8032",
    "working_dir" = "hdfs://<namenode-host>:<port>/tmp/starrocks",
    "broker" = "<broker-name>",
    "broker.username" = "<username>",
    "broker.password" = "<password>"
);

リソースプロパティ

プロパティ必須説明
typeはいspark に設定します。
spark.masterはいyarn に設定します。
spark.submit.deployModeはいSpark プログラムのデプロイモード。有効な値:clusterclient
spark.hadoop.fs.defaultFSはいspark.masteryarn の場合に必須です。
spark.hadoop.yarn.resourcemanager.addressいいえYARN ResourceManager のアドレス。host:port フォーマットです。
spark.hadoop.yarn.resourcemanager.ha.enabledいいえResourceManager で高可用性を有効にするには true に設定します。デフォルト値:true
spark.hadoop.yarn.resourcemanager.ha.rm-idsいいえResourceManager の論理 ID (HA 用)。
spark.hadoop.yarn.resourcemanager.hostname.rm-idいいえ各論理 ID に対応するホスト名。これか spark.hadoop.yarn.resourcemanager.address.rm-id のいずれかを設定します。
spark.hadoop.yarn.resourcemanager.address.rm-idいいえ各論理 ID に対応するアドレス (host:port)。これか spark.hadoop.yarn.resourcemanager.hostname.rm-id のいずれかを設定します。
working_dirはい (ETL 用)Spark ETL リソースがステージングされる HDFS ディレクトリ。例:hdfs://host:port/tmp/starrocks
brokerはい (ETL 用)使用する Broker の名前。ALTER SYSTEM ADD BROKER を実行して先に追加します。
broker.property_keyいいえBroker が ETL 中間ファイルを読み取るために使用する認証プロパティ。

Spark 設定オプションの完全なリストについては、Spark 設定ドキュメントをご参照ください。

Yarn クラスターモード (標準):

CREATE EXTERNAL RESOURCE "spark0"
PROPERTIES
(
    "type" = "spark",
    "spark.master" = "yarn",
    "spark.submit.deployMode" = "cluster",
    "spark.jars" = "xxx.jar,yyy.jar",
    "spark.files" = "/tmp/aaa,/tmp/bbb",
    "spark.executor.memory" = "1g",
    "spark.yarn.queue" = "queue0",
    "spark.hadoop.yarn.resourcemanager.address" = "resourcemanager_host:8032",
    "spark.hadoop.fs.defaultFS" = "hdfs://namenode_host:9000",
    "working_dir" = "hdfs://namenode_host:9000/tmp/starrocks",
    "broker" = "broker0",
    "broker.username" = "user0",
    "broker.password" = "password0"
);

Yarn クラスターモード (YARN 高可用性):

CREATE EXTERNAL RESOURCE "spark1"
PROPERTIES
(
    "type" = "spark",
    "spark.master" = "yarn",
    "spark.submit.deployMode" = "cluster",
    "spark.hadoop.yarn.resourcemanager.ha.enabled" = "true",
    "spark.hadoop.yarn.resourcemanager.ha.rm-ids" = "rm1,rm2",
    "spark.hadoop.yarn.resourcemanager.hostname.rm1" = "host1",
    "spark.hadoop.yarn.resourcemanager.hostname.rm2" = "host2",
    "spark.hadoop.fs.defaultFS" = "hdfs://namenode_host:9000",
    "working_dir" = "hdfs://namenode_host:9000/tmp/starrocks",
    "broker" = "broker1"
);

HDFS 高可用性モード:

CREATE EXTERNAL RESOURCE "spark2"
PROPERTIES
(
    "type" = "spark",
    "spark.master" = "yarn",
    "spark.hadoop.yarn.resourcemanager.address" = "resourcemanager_host:8032",
    "spark.hadoop.fs.defaultFS" = "hdfs://myha",
    "spark.hadoop.dfs.nameservices" = "myha",
    "spark.hadoop.dfs.ha.namenodes.myha" = "mynamenode1,mynamenode2",
    "spark.hadoop.dfs.namenode.rpc-address.myha.mynamenode1" = "nn1_host:rpc_port",
    "spark.hadoop.dfs.namenode.rpc-address.myha.mynamenode2" = "nn2_host:rpc_port",
    "spark.hadoop.dfs.client.failover.proxy.provider" = "org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider",
    "working_dir" = "hdfs://myha/tmp/starrocks",
    "broker" = "broker2",
    "broker.dfs.nameservices" = "myha",
    "broker.dfs.ha.namenodes.myha" = "mynamenode1,mynamenode2",
    "broker.dfs.namenode.rpc-address.myha.mynamenode1" = "nn1_host:rpc_port",
    "broker.dfs.namenode.rpc-address.myha.mynamenode2" = "nn2_host:rpc_port",
    "broker.dfs.client.failover.proxy.provider" = "org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider"
);

リソースの管理

-- すべてのリソースをリスト表示
SHOW RESOURCES;
SHOW PROC "/resources";

-- USAGE-PRIV の付与または取り消し
GRANT USAGE_PRIV ON RESOURCE "spark0" TO "user0"@"%";
GRANT USAGE_PRIV ON RESOURCE "spark0" TO ROLE "role0";
GRANT USAGE_PRIV ON RESOURCE * TO "user0"@"%";
GRANT USAGE_PRIV ON RESOURCE * TO ROLE "role0";
REVOKE USAGE_PRIV ON RESOURCE "spark0" FROM "user0"@"%";
REVOKE USAGE_PRIV ON RESOURCE "spark0" FROM ROLE "role0";
通常のアカウントでは、SHOW RESOURCES はアカウントが USAGE-PRIV を持つリソースのみを表示します。root および admin アカウントはすべてのリソースを表示します。

リソースの削除

DROP RESOURCE "resource_name";

ステップ 2:Spark クライアントの設定

フロントエンドノードは spark-submit を使用して Spark Load ジョブを送信するため、Spark クライアントをフロントエンドノードのホストにインストールする必要があります。

Spark 2.4.5 またはそれ以降の 2.x バージョンをダウンロードしてフロントエンドノードに配置します。次に、フロントエンドノードの設定ファイル (fe.conf) で以下のパラメーターを設定します。

  1. Spark ホームディレクトリの設定spark_home_default_dir を、Spark クライアントを配置したディレクトリに設定します。デフォルト値は、フロントエンドノードのルートディレクトリ下の lib/spark2x です。このパラメーターは空白にできません。

  2. Spark 依存関係のパッケージ化:Spark クライアントの jars/ フォルダにあるすべての JAR ファイルを ZIP ファイルにパッケージ化します。spark_resource_path をその ZIP ファイルのパスに設定します。空白のままにすると、フロントエンドノードはルートディレクトリ内の lib/spark2x/jars/spark-2x.zip を探します。見つからない場合はエラーが返されます。Spark Load ジョブが送信されると、依存関係パッケージは working_dir/{cluster_id}/ のリモートステージング場所にアップロードされます。ステージング場所は --spark-repository--{resource-name} のフォーマットで命名されます。ディレクトリ構造の例:

    ---spark-repository--spark0/
       |---archive-1.0.0/
       |   |---lib-990325d2c0d1d5e45bf675e54e44fb16-spark-dpp-1.0.0-jar-with-dependencies.jar
       |   |---lib-7670c29daf535efe3c9b923f778f61fc-spark-2x.zip
       |---archive-1.1.0/
       |   |---lib-64d5696f99c379af2bee28c1c84271d5-spark-dpp-1.1.0-jar-with-dependencies.jar
       |   |---lib-1bbb74bb6b264a270bc7fca3e964160f-spark-2x.zip

    デフォルトの依存関係パッケージ名は spark-2x.zip です。フロントエンドノードは、動的パーティションプルーニング (DPP) の依存関係もアップロードします。一度アップロードされると、これらのパッケージは後続のジョブで再利用されます。

ステップ 3:YARN クライアントの設定

フロントエンドノードは YARN コマンドを実行してアプリケーションのステータスを確認したり、アプリケーションを停止したりするため、YARN クライアントもフロントエンドノードにインストールする必要があります。

Hadoop 2.5.2 またはそれ以降の 2.x バージョンをダウンロードし、fe.conf で以下のパラメーターを設定します。

  1. YARN 実行可能ファイルのパスの設定yarn_client_path を YARN バイナリファイルのパスに設定します。デフォルトは、フロントエンドノードのルートディレクトリ下の lib/yarn-client/hadoop/bin/yarn です。

  2. (オプション) YARN 設定ディレクトリの設定:フロントエンドノードがアプリケーションのステータスを確認したり、アプリケーションを停止したりすると、デフォルトで lib/yarn-config/ 内に core-site.xmlyarn-site.xml が生成されます。この場所を変更するには、yarn_config_dir を設定します。

データのインポート

インポートジョブの作成

サポートされているデータソース:CSV ファイルと Hive テーブル。

LOAD LABEL db_name.label_name
    (data_desc, ...)
WITH RESOURCE resource_name
[resource_properties]
[PROPERTIES (key1=value1, ...)]

ここで、data_desc は以下のいずれかです。

-- HDFS ファイルから
DATA INFILE ("file_path", ...)
[NEGATIVE]
INTO TABLE tbl_name
[PARTITION (p1, p2)]
[COLUMNS TERMINATED BY separator]
[(col1, ...)]
[COLUMNS FROM PATH AS (col2, ...)]
[SET (k1=f1(xx), k2=f2(xx))]
[WHERE predicate]

-- Hive テーブルから
DATA FROM TABLE hive_external_tbl
[NEGATIVE]
INTO TABLE tbl_name
[PARTITION (p1, p2)]
[SET (k1=f1(xx), k2=f2(xx))]
[WHERE predicate]

完全な構文については、HELP SPARK LOAD を実行してください。

主要なパラメーター

パラメーター説明
labelデータベース内でインポートジョブを一意に識別する識別子。Broker Load と同じ仕様です。
データ記述データソースとして CSV ファイルと Hive テーブルをサポートします。その他の仕様は Broker Load と同じです。
ジョブプロパティBroker Load の opt_properties と同じです。
Spark リソースプロパティジョブごとの Spark リソースのオーバーライド。このジョブにのみ有効で、クラスターレベルのリソース設定は変更されません。

例 1:HDFS CSV ファイルからのインポート

LOAD LABEL db1.label1
(
    DATA INFILE("hdfs://emr-header-1.cluster-xxx:9000/user/starRocks/test/ml/file1")
    INTO TABLE tbl1
    COLUMNS TERMINATED BY ","
    (tmp_c1, tmp_c2)
    SET
    (
        id=tmp_c2,
        name=tmp_c1
    ),
    DATA INFILE("hdfs://emr-header-1.cluster-xxx:9000/user/starRocks/test/ml/file2")
    INTO TABLE tbl2
    COLUMNS TERMINATED BY ","
    (col1, col2)
    WHERE col1 > 1
)
WITH RESOURCE 'spark0'
(
    "spark.executor.memory" = "2g",
    "spark.shuffle.compress" = "true"
)
PROPERTIES
(
    "timeout" = "3600"
);

例 2:グローバル辞書を使用した Hive テーブルからのインポート

このアプローチは、ターゲットの StarRocks テーブルにビットマップ集約列がある場合に使用します。bitmap_dict() 関数は、グローバル辞書を介して生の Hive 列の値をエンコードされた整数にマッピングします。

  1. Hive 外部リソースを作成します。

    CREATE EXTERNAL RESOURCE hive0
    PROPERTIES
    (
        "type" = "hive",
        "hive.metastore.uris" = "thrift://emr-header-1.cluster-xxx:9083"
    );
  2. 外部 Hive テーブルを作成します。

    CREATE EXTERNAL TABLE hive_t1
    (
        k1 INT,
        K2 SMALLINT,
        k3 VARCHAR(50),
        uuid VARCHAR(100)
    )
    ENGINE=hive
    PROPERTIES
    (
        "resource" = "hive0",
        "database" = "tmp",
        "table" = "t1"
    );
  3. インポートジョブを送信します。

    LOAD LABEL db1.label1
    (
        DATA FROM TABLE hive_t1
        INTO TABLE tbl1
        SET
        (
            uuid=bitmap_dict(uuid)
        )
    )
    WITH RESOURCE 'spark0'
    (
        "spark.executor.memory" = "2g",
        "spark.shuffle.compress" = "true"
    )
    PROPERTIES
    (
        "timeout" = "3600"
    );

インポートジョブの確認

SHOW LOAD ORDER BY createtime DESC LIMIT 1\G

出力例:

*************************** 1. row ***************************
       JobId: 76391
       Label: label1
       State: FINISHED
    Progress: ETL:100%; LOAD:100%
        Type: SPARK
     EtlInfo: unselected.rows=4; dpp.abnorm.ALL=15; dpp.norm.ALL=28133376
    TaskInfo: cluster:cluster0; timeout(s):10800; max_filter_ratio:5.0E-5
    ErrorMsg: N/A
  CreateTime: 2019-07-27 11:46:42
EtlStartTime: 2019-07-27 11:46:44
EtlFinishTime: 2019-07-27 11:49:44
LoadStartTime: 2019-07-27 11:49:44
LoadFinishTime: 2019-07-27 11:50:16
         URL: http://1.1.*.*:8089/proxy/application_1586619723848_0035/
  JobDetails: {"ScannedRows":28133395,"TaskNumber":1,"FileNumber":1,"FileSize":200000}

出力フィールド

フィールド説明
Stateジョブのステータス。PENDING → ETL → LOADING → FINISHED (失敗した場合は CANCELLED) に遷移します。
ProgressETL と LOAD の進捗率。LOAD の進捗率は、すべてのレプリカでインポートされたタブレット数 / 総タブレット数 × 100% として計算されます。ジョブはインポートが有効になる前に 99% に達し、その後 100% にジャンプします。進捗は必ずしも線形ではありません。
Typeインポートタイプ。Spark Load ジョブの場合は SPARK と表示されます。
CreateTimeインポートジョブが作成された時間。
EtlStartTimeジョブが ETL 状態に入った時間。
EtlFinishTimeジョブが ETL 状態を抜けた時間。
LoadStartTimeジョブが LOADING 状態に入った時間。
LoadFinishTimeジョブが完了した時間。
JobDetailsスキャンされた行数、タスク数、ファイル数、総データサイズなどの詳細。例:{"ScannedRows":139264,"TaskNumber":1,"FileNumber":1,"FileSize":940754064}
URLSpark アプリケーションの Web UI の URL。ブラウザで開くとジョブの詳細が表示されます。

完全なパラメーターリファレンスについては、「Broker Load」をご参照ください。

完全な構文については、HELP SHOW LOAD を実行してください。

ジョブログの表示

Spark Load のログは、フロントエンドノードのルートディレクトリ下の log/spark_launcher_log/ に保存され、spark-launcher-{load-job-id}-{label}.log のフォーマットで命名されます。ログはデフォルトで 3 日間保持され、関連するインポートメタデータがパージされるときに削除されます。

インポートジョブのキャンセル

FINISHED または CANCELLED 状態でないジョブをキャンセルします。

CANCEL LOAD WHERE LABEL = "label1";

完全な構文については、HELP CANCEL LOAD を実行してください。

システム設定

fe.conf 内のこれらのパラメーターは、システムレベルで Spark Load の動作を制御します。

パラメーターデフォルト説明
enable_spark_loadfalseSpark Load と外部リソースの作成を有効にします。有効にするには true に設定します。
spark_load_default_timeout_second259200 (3 日)インポートジョブのデフォルトのタイムアウト (秒単位)。
spark_home_default_dirfe/lib/spark2xSpark クライアントが保存されているディレクトリ。
spark_resource_path(空)パッケージ化された Spark 依存関係 ZIP ファイルへのパス。
spark_launcher_log_dirfe/log/spark-launcher-logSpark クライアントの送信ログが保存されるディレクトリ。
yarn_client_pathfe/lib/yarn-client/hadoop/bin/yarnYARN クライアントバイナリへのパス。
yarn_config_dirfe/lib/yarn-configYARN コマンドの設定ファイルが生成されるディレクトリ。

ベストプラクティス

HDFS からの数十 GB から TB 範囲のインポートには Spark Load を使用してください。より小さなデータセットの場合は、Stream Load または Broker Load の方が適しています。これらは設定のオーバーヘッドが低く、処理時間が短いためです。

完全なエンドツーエンドのサンプルコードについては、GitHub の 03_sparkLoad2StarRocks.md をご参照ください。

次のステップ

  • Broker Load — より小さなデータセット向けの代替インポートメソッド

  • Resource Management — StarRocks で外部の Spark リソースを管理する