StarRocks は、さまざまなビジネスシナリオに対応する複数のデータモデルをサポートしており、すべてのデータは特定のモデルに従って構成する必要があります。このトピックでは、さまざまなインポート方法に関する基本概念、原則、システム構成、ユースケース、ベストプラクティス、FAQ について説明します。
背景情報
データインポートは、特定のデータモデルに基づいて生データをクレンジング、変換し、クエリのために StarRocks にロードします。StarRocks は複数のインポート方法を提供しており、データ量、インポート頻度、その他のビジネス要件に基づいて選択できます。
次の図は、StarRocks のインポート方法とさまざまなデータソースとの関係を示しています。
データソースに基づいて、さまざまなインポート方法を選択できます。
-
オフラインデータインポート:ソースデータが Hive または HDFS にある場合は、Broker Load を使用します。多くのデータテーブルが関わる複雑なインポートには、 を使用できます。この方法は Broker Load よりもパフォーマンスが低いですが、データ移行を回避できます。単一テーブルのデータ量が非常に大きい場合や、正確な重複排除のためのグローバルデータ辞書として使用する場合は、Spark Load の使用を検討してください。
-
リアルタイムデータインポート:ログデータとデータベースのバイナリロギング (binlog) を Kafka に同期した後、Routine Load を使用してデータを StarRocks にインポートできます。インポートプロセスに複雑なテーブル結合や抽出、変換、ロード (ETL) の前処理が含まれる場合は、Flink (Flink コネクタ) を使用して最初にデータを処理できます。その後、Stream Load を使用してデータを StarRocks に書き込むことができます。
-
プログラムから StarRocks へデータを書き込むには、Stream Load を使用できます。例については、Stream Load ドキュメントの Java または Python のデモをご参照ください。
-
テキストファイルのインポート:Stream Load を使用できます。
-
MySQL データのインポート:MySQL 外部テーブルを使用できます。データをインポートするには、
insert into new_table select * from external_tableコマンドを実行します。 -
StarRocks 内部でのインポート:Insert Into メソッドを使用できます。このメソッドは、外部スケジューラと連携して単純な ETL 処理を実行します。
このトピックの画像と一部のコンテンツは、オープンソース StarRocks のドキュメント「データロードの概要」からの引用です。
注意事項
データを StarRocks にインポートする場合、通常はプログラムを使用して接続を確立します。次の点にご注意ください。
-
適切なインポート方法の選択:データ量、インポート頻度、データソースの場所に基づいてインポート方法を選択します。
例えば、ソースデータが HDFS にある場合、Broker Load を使用できます。
-
インポート方法のプロトコルの決定:Broker Load を選択した場合、外部システムは MySQL プロトコルを使用して定期的にインポートジョブを送信し、確認できる必要があります。
-
インポート方法の種類の決定:インポート方法には同期と非同期があります。非同期インポート方法の場合、インポートジョブを送信した後、コマンドを実行してジョブのステータスを表示する必要があります。このコマンドの結果で、インポートが成功したかどうかがわかります。
-
ラベル生成ポリシーの作成:ポリシーは、各バッチのデータに対して各ラベルが一意で固定されていることを保証する必要があります。
-
Exactly-Once セマンティクスの保証:外部システムは At-Least-Once のデータインポートを保証する必要があります。StarRocks のラベルメカニズムは At-Most-Once のデータインポートを保証します。この 2 つのメカニズムを組み合わせることで、データインポートプロセス全体で Exactly-Once セマンティクスが保証されます。
用語
|
用語 |
説明 |
|
インポートジョブ |
ユーザーが送信したソースデータを読み取り、クレンジングと変換を行った後、StarRocks システムにロードします。インポートが完了すると、データはクエリ可能になります。 |
|
ラベル |
インポートジョブを識別します。すべてのインポートジョブにはラベルがあります。 ラベルはユーザーが指定するか、システムが生成し、データベース内で一意です。ラベルは 1 つの成功したインポートジョブにのみ使用できます。特定のラベルを持つインポートジョブが成功した後、そのラベルを再利用して別のインポートジョブを送信することはできません。インポートジョブが失敗した場合、ラベルは再利用できます。このメカニズムにより、At-Most-Once のインポートセマンティクスが保証されます。 |
|
原子性 |
StarRocks のすべてのインポート方法は原子性を提供します。これは、1 つのインポートジョブに対して、すべての有効なデータが正常にインポートされるか、まったくインポートされないかのどちらかであることを意味します。部分的なインポートは発生しません。ここでの有効なデータには、型変換エラーなどの品質問題のために除外されたデータは含まれません。データ品質問題の詳細については、「データインポートに関する FAQ」をご参照ください。 |
|
MySQL および HTTP プロトコル |
StarRocks は、ジョブを送信するための 2 つのアクセスプロトコルインターフェイス (MySQL プロトコルと HTTP プロトコル) を提供します。 |
|
Broker Load |
デプロイされたブローカープログラムを介して HDFS などの外部ソースからデータを読み取り、StarRocks にインポートします。ブローカープロセスは、独自のコンピューティングリソースを使用してデータを前処理します。 |
|
Spark Load |
外部の Spark リソースを使用してデータを前処理し、StarRocks がインポートのために読み取る中間ファイルを生成します。Spark Load は非同期です。MySQL プロトコルを使用してインポートジョブを作成し、 |
|
FE |
Frontend (FE)。StarRocks のメタデータおよびスケジューリングノードで、インポート実行計画の生成とインポートタスクのスケジューリングを担当します。 |
|
BE |
Backend (BE)。StarRocks のコンピューティングおよびストレージノードで、インポート中のデータの ETL とストレージを担当します。 |
|
タブレット |
StarRocks テーブルの論理的なシャードです。テーブルは、パーティショニングとバケット化のルールに従って複数のタブレットに分割できます。詳細については、「データパーティション」をご参照ください。 |
基本原則
インポートの実行フローは次の図のとおりです。
インポートジョブは、次の 5 つのステージで構成されます。
|
ステージ |
説明 |
|
PENDING |
オプション。このステージでは、インポートジョブが送信され、FE が実行をスケジュールするのを待機します。 Broker Load と Spark Load にはこのステージが含まれます。 |
|
ETL |
オプション。このステージでは、データのクレンジング、パーティショニング、ソート、集計などの前処理を実行します。 Spark Load にはこのステージが含まれます。外部のコンピューティングリソースである Spark を使用して ETL プロセスを完了します。 |
|
LOADING |
このステージでは、データが最初にクレンジングされ、変換された後、処理のために BE に送信されます。すべてのデータがロードされると、ジョブはデータが有効になるのを待つ状態になります。この時点で、インポートジョブのステータスはまだ LOADING です。 |
|
FINISHED |
インポートジョブの対象となるすべてのデータが有効になると、ジョブのステータスは FINISHED に変わります。FINISHED 状態のジョブのデータはクエリできます。FINISHED は、成功したインポートジョブの最終状態です。 |
|
CANCELLED |
ジョブのステータスが FINISHED に変わる前に、いつでもキャンセルでき、CANCELLED 状態になります。例えば、手動でキャンセルしたり、インポート中にエラーが発生したりすることがあります。CANCELLED もインポートジョブの最終状態です。 |
データインポートのフォーマットを次の表に示します。
|
タイプ |
説明 |
|
整数型 |
TINYINT、SMALLINT、INT、BIGINT、および LARGEINT。例:1、1000、1234。 |
|
浮動小数点型 |
FLOAT、DOUBLE、および DECIMAL。例:1.1、0.23、0.356。 |
|
日付型 |
DATE および DATETIME。例:2017-10-03、2017-06-13 12:34:03。 |
|
文字列型 |
CHAR および VARCHAR。例:I am a student、a。 |
インポート方法
StarRocks は、HDFS、Kafka、ローカルファイルなどのさまざまなデータソース向けに 5 つのインポート方法を提供します。これらの方法は、同期的または非同期的です。
すべてのインポート方法は CSV データ形式をサポートしています。Broker Load は、Parquet および ORC データ形式もサポートしています。
インポート方法の紹介
|
インポート方法 |
説明 |
インポートタイプ |
|
Broker Load |
ブローカープロセスを通じて外部データソースを読み取り、MySQL プロトコルを介して StarRocks にインポートジョブを作成します。ジョブは非同期で実行されます。 HDFS などのブローカーがアクセス可能なストレージシステム内のデータに適しており、データ量は数十から数百ギガバイトです。詳細については、「Broker Load」をご参照ください。 |
非同期インポート |
|
Spark Load |
外部の Spark リソースを使用してインポートされたデータを前処理し、大規模なデータセットのパフォーマンスを向上させ、StarRocks クラスターのリソース使用量を削減します。この非同期メソッドでは、MySQL プロトコルを介してジョブを作成し、 Spark Load は、大量のデータ (最大テラバイトレベル) を初めて StarRocks に移行するのに適しています。ソースデータは、Spark がアクセスできるストレージシステム (HDFS など) にある必要があります。詳細については、「Spark Load」をご参照ください。 |
非同期インポート |
|
Stream Load |
HTTP プロトコルを介してローカルファイルまたはデータストリームを StarRocks にインポートする同期メソッドです。インポート結果はレスポンスで直接返されます。 ローカルファイルのインポートやプログラムからのストリーミングデータのインポートに適しています。詳細については、「Stream Load」をご参照ください。 |
同期インポート |
|
Routine Load |
永続的なスレッドを作成することで、指定されたソースからデータを自動的にインポートします。MySQL プロトコルを通じてルーチンロードジョブを送信し、Kafka などのソースからデータを継続的に読み取ってインポートします。詳細については、「Routine Load」をご参照ください。 |
非同期インポート |
|
Insert Into |
MySQL の |
同期インポート |
インポートタイプ
外部プログラムが StarRocks のデータインポート機能を使用する場合、まずインポート方法のタイプを決定し、次に接続ロジックを定義する必要があります。
-
同期インポート
同期インポートでは、StarRocks はタスクをすぐに実行し、インポートが成功したかどうかを示す結果を返します。
手順:
-
ユーザー (外部システム) がインポートタスクを作成します。
-
StarRocks はインポート結果を返します。
-
ユーザー (外部システム) がインポート結果を確認します。インポートが失敗した場合、ユーザーは再度インポートタスクを作成できます。
-
-
非同期インポート
非同期インポートでは、StarRocks は作成成功のメッセージをすぐに返しますが、データはまだインポートされていません。コマンドを実行してジョブのステータスをポーリングする必要があります。タスクの作成が失敗した場合、失敗情報に基づいて再試行できます。
手順:
-
ユーザー (外部システム) がインポートタスクを作成します。
-
StarRocks はタスク作成の結果を返します。
-
ユーザー (外部システム) がタスク作成の結果を確認します。タスクが正常に作成された場合は、手順 4 に進みます。それ以外の場合は、手順 1 に戻り、インポートタスクの作成を再試行します。
-
ユーザー (外部システム) は、ステータスが FINISHED または CANCELLED になるまでタスクのステータスをポーリングします。
-
シナリオ
|
シナリオ |
説明 |
|
HDFS インポート |
ソースデータが HDFS に保存されており、データ量が数十から数百ギガバイトの場合、Broker Load メソッドを使用してデータを StarRocks にインポートできます。デプロイされたブローカープロセスは、HDFS データソースにアクセス可能である必要があります。インポートジョブは非同期で実行されます。 ソースデータが HDFS に保存されており、データ量がテラバイトレベルに達する場合、Spark Load メソッドを使用してデータを StarRocks にインポートできます。デプロイされた Spark プロセスは、HDFS データソースにアクセス可能である必要があります。インポートジョブは非同期で実行されます。 他の外部データソースについても、ブローカーまたは Spark プロセスが対応するデータソースから読み取ることができれば、Broker Load または Spark Load を使用してデータをインポートできます。 |
|
ローカルファイルインポート |
10 GB 未満のローカルファイルには、Stream Load を使用します。HTTP プロトコルを介してインポートジョブを作成します。ジョブは同期的に実行され、結果を直接返します。 |
|
Kafka インポート |
Kafka などのストリーミングソースからのリアルタイムデータには、Routine Load を使用します。MySQL プロトコルを介してルーチンロードジョブを作成すると、StarRocks は継続的にデータを読み取り、インポートします。 |
|
Insert Into インポート |
手動テストや一時的なデータ処理には、
|
メモリ制限
パラメータを設定して、インポートジョブごとのメモリ使用量を制限し、メモリ不足 (OOM) エラーを防ぎます。メモリを制限する方法はインポート方法によって異なります。詳細については、各メソッドのドキュメントをご参照ください。
インポートジョブは通常、複数の BE に分散されます。メモリ制限パラメータは、クラスター全体ではなく、単一の BE 上での 1 つのインポートジョブのメモリ使用量を制限します。各 BE には、すべてのインポートジョブに対する合計メモリ制限もあります。詳細については、「一般的なシステム構成」をご参照ください。
メモリ制限が小さいと、頻繁なディスク書き込みが発生してインポート効率が低下する可能性があり、制限が大きいと、同時実行性が高い場合に OOM エラーのリスクがあります。ワークロード要件に基づいてメモリパラメータを設定してください。
一般的なシステム構成
FE の構成
次の FE パラメータを fe.conf ファイルで設定します。
|
パラメータ |
説明 |
|
max_load_timeout_second |
インポートジョブの最大および最小タイムアウト期間 (秒単位)。デフォルトの最大タイムアウトは 3 日、デフォルトの最小タイムアウトは 1 秒です。インポートジョブに設定するカスタムタイムアウトは、この範囲を超えることはできません。このパラメータは、すべてのタイプのインポートタスクに適用されます。 |
|
min_load_timeout_second |
|
|
desired_max_waiting_jobs |
待機キューが保持できるインポートタスクの最大数。デフォルト値は 100 です。 例えば、FE 上の PENDING 状態 (実行待ち) のインポートタスクの数がこの値に達すると、新しいインポートリクエストは拒否されます。この構成は、非同期で実行されるインポートにのみ適用されます。待機中の非同期インポートタスクの数が制限に達すると、後続のインポートジョブ作成リクエストは拒否されます。 |
|
max_running_txn_num_per_db |
各データベースで実行中のインポートタスクの最大数。これはすべてのインポートタイプにわたってカウントされます。デフォルト値は 100 です。 データベースで実行中のインポートタスクの数が最大値を超えると、後続のインポートタスクは実行されません。同期ジョブの場合、ジョブは拒否されます。非同期ジョブの場合、ジョブはキューで待機します。 |
|
label_keep_max_second |
インポートタスクレコードの保持期間。 完了した (FINISHED または CANCELLED) インポートタスクのレコードは、このパラメータで指定された期間、StarRocks システムに保持されます。デフォルト値は 3 日です。このパラメータは、すべてのタイプのインポートタスクに適用されます。 |
BE の構成
次の BE パラメータを be.conf ファイルで設定します。
|
パラメータ |
説明 |
|
push_write_mbytes_per_sec |
BE 上の単一タブレットの書き込み速度制限。デフォルト値は 10 で、10 MB/s を意味します。 BE 上の単一タブレットの最大書き込み速度は、スキーマとシステムに応じて、通常 10 MB/s から 30 MB/s の範囲です。このパラメータを調整して、インポート速度を制御できます。 |
|
write_buffer_size |
データインポート中、データはまず BE 上のメモリブロックに書き込まれます。このメモリブロックがしきい値に達すると、ディスクに書き込まれます。デフォルト値は 100 MB です。 しきい値が小さいと、BE 上に多くの小さなファイルが作成される可能性があります。このしきい値を増やすことで、ファイル数を減らすことができます。ただし、しきい値が大きいと RPC タイムアウトが発生する可能性があります。詳細については、tablet_writer_rpc_timeout_sec パラメータをご参照ください。 |
|
tablet_writer_rpc_timeout_sec |
インポートプロセス中にバッチ (1024 行) を送信するための RPC タイムアウト。デフォルトは 600 秒です。 この RPC には、複数のタブレットメモリブロックをディスクに書き込むことが含まれる場合があります。したがって、ディスク書き込みが原因で RPC タイムアウトが発生する可能性があります。タイムアウトを調整して、 |
|
streaming_load_rpc_max_alive_time_sec |
インポートプロセス中、StarRocks は各タブレットのライターを起動してデータを受信し、書き込みます。このパラメータは、ライターの待機タイムアウトを指定します。デフォルトは 600 秒です。 指定された時間内にライターがデータを受信しない場合、ライターは自動的に破棄されます。システムの処理速度が遅い場合、ライターは次のデータバッチを長時間受信できず、「 |
|
load_process_max_memory_limit_percent |
これらのパラメータは、インポートタスクが使用できるメモリの上限を、それぞれ絶対値と割合で指定します。これらは、単一の BE 上でインポートタスクに使用できる合計メモリを制限します。システムは、load_process_max_memory_limit_percent から算出されるメモリ上限 (バイト単位) と、load_process_max_memory_limit_bytes で指定されるメモリ上限を比較し、小さい方を最終的な制限として使用します。
|
|
load_process_max_memory_limit_bytes |