Spark Load は、StarRocks にデータをインポートする前に、外部の Spark クラスターにデータの前処理をオフロードします。これにより、StarRocks クラスターの計算負荷が軽減され、初回移行や TB レベルの大規模なインポートに推奨されるアプローチです。
ユーザーアカウントは、インポートジョブを送信する前に、Spark リソースに対する USAGE-PRIV 権限を持っている必要があります。
仕組み
Spark Load は非同期のインポートメソッドです。MySQL プロトコルを使用してジョブを送信し、SHOW LOAD で結果を確認します。
以下の図は、そのワークフローを示しています。

フロントエンドノードに Spark Load ジョブを送信します。
フロントエンドノードは、抽出、変換、ロード (ETL) ジョブをスケジュールし、Spark クラスターに送信します。
Spark クラスターは ETL ジョブを実行します。ビットマップグローバル辞書を構築し、データのパーティション分割、ソート、集約を行います。
ETL ジョブが完了すると、フロントエンドノードは各パーティションの前処理済みデータディレクトリを特定し、バックエンドノードにプッシュジョブを実行するようスケジュールします。
バックエンドノードは Broker を使用して Hadoop 分散ファイルシステム (HDFS) からデータを読み取り、StarRocks が内部で保存するフォーマットに変換します。
フロントエンドノードは新しい StarRocks バージョンを公開し、インポートジョブを完了としてマークします。
基本概念
Spark ETL:インポート前にデータに対して ETL 操作を実行する Spark プログラムです。グローバル辞書の作成、データのパーティション分割、ソート、集約を処理します。
Broker:ファイルシステムインターフェイスをラップするステートレスなプロセスで、StarRocks が HDFS などのリモートストレージシステムからファイルを読み取れるようにします。
グローバル辞書:生の値とエンコードされた整数をマッピングするデータ構造です。グローバル辞書は、インポート前にビットマップ列を事前計算するために使用されます。StarRocks のビットマップ列は roaring bitmap を使用しており、整数の入力が必要です。
グローバル辞書のワークフロー
StarRocks のビットマップ列には整数の入力が必要です。グローバル辞書がこの変換を処理します。
ソースデータを読み取り、一時的な Hive テーブル (
hive-table) に保存します。hive-tableから値を重複排除し、distinct-value-tableという名前のテーブルに格納します。辞書テーブル (
dict-table) を作成します。このテーブルには、生の値用の列とエンコードされた整数用の列が 1 つずつあります。distinct-value-tableとdict-tableを LEFT JOIN し、ウィンドウ関数を適用して新しい生の値をエンコードし、結果をdict-tableに書き戻します。dict-tableとhive-tableを JOIN して、生の値をエンコードされた整数に置き換えます。エンコードされた
hive-tableデータを後続の ETL ステップに渡し、StarRocks にインポートします。
グローバル辞書は、Hive テーブルからインポートする場合にのみサポートされます。
データの前処理
グローバル辞書が構築された後 (該当する場合)、Spark はデータを前処理します。
HDFS ファイルまたは Hive テーブルからデータを読み取ります。
フィールドマッピングと式ベースの計算を適用します。パーティション情報に基づいて
bucket-idフィールドを生成します。StarRocks テーブルのロールアップメタデータからロールアップツリーを構築します。
ロールアップツリーを走査し、データをレイヤーごとに集約します。各レイヤーは前のレイヤーから計算されます。
bucket-idに基づいて集約データをバケットに分散し、HDFS に書き込みます。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 プログラムのデプロイモード。有効な値:cluster、client。 |
spark.hadoop.fs.defaultFS | はい | spark.master が yarn の場合に必須です。 |
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) で以下のパラメーターを設定します。
Spark ホームディレクトリの設定:
spark_home_default_dirを、Spark クライアントを配置したディレクトリに設定します。デフォルト値は、フロントエンドノードのルートディレクトリ下のlib/spark2xです。このパラメーターは空白にできません。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 で以下のパラメーターを設定します。
YARN 実行可能ファイルのパスの設定:
yarn_client_pathを YARN バイナリファイルのパスに設定します。デフォルトは、フロントエンドノードのルートディレクトリ下のlib/yarn-client/hadoop/bin/yarnです。(オプション) YARN 設定ディレクトリの設定:フロントエンドノードがアプリケーションのステータスを確認したり、アプリケーションを停止したりすると、デフォルトで
lib/yarn-config/内にcore-site.xmlとyarn-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 列の値をエンコードされた整数にマッピングします。
Hive 外部リソースを作成します。
CREATE EXTERNAL RESOURCE hive0 PROPERTIES ( "type" = "hive", "hive.metastore.uris" = "thrift://emr-header-1.cluster-xxx:9083" );外部 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" );インポートジョブを送信します。
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) に遷移します。 |
Progress | ETL と 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}。 |
URL | Spark アプリケーションの 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_load | false | Spark Load と外部リソースの作成を有効にします。有効にするには true に設定します。 |
spark_load_default_timeout_second | 259200 (3 日) | インポートジョブのデフォルトのタイムアウト (秒単位)。 |
spark_home_default_dir | fe/lib/spark2x | Spark クライアントが保存されているディレクトリ。 |
spark_resource_path | (空) | パッケージ化された Spark 依存関係 ZIP ファイルへのパス。 |
spark_launcher_log_dir | fe/log/spark-launcher-log | Spark クライアントの送信ログが保存されるディレクトリ。 |
yarn_client_path | fe/lib/yarn-client/hadoop/bin/yarn | YARN クライアントバイナリへのパス。 |
yarn_config_dir | fe/lib/yarn-config | YARN コマンドの設定ファイルが生成されるディレクトリ。 |
ベストプラクティス
HDFS からの数十 GB から TB 範囲のインポートには Spark Load を使用してください。より小さなデータセットの場合は、Stream Load または Broker Load の方が適しています。これらは設定のオーバーヘッドが低く、処理時間が短いためです。
完全なエンドツーエンドのサンプルコードについては、GitHub の 03_sparkLoad2StarRocks.md をご参照ください。
次のステップ
Broker Load — より小さなデータセット向けの代替インポートメソッド
Resource Management — StarRocks で外部の Spark リソースを管理する