ストリームロードを使用して、ローカルファイルまたはデータストリームから StarRocks にデータをロードします。
背景情報
ストリームロードは、HTTP リクエストを使用してローカルファイルまたはデータストリームから StarRocks にデータをロードする同期ロード方式です。操作は同期的であるため、HTTP 応答は即座にロード結果を返し、その成功を確認できます。ストリームロードは CSV と JSON 形式をサポートしており、各ロードは 10 GB 未満である必要があります。
ロードタスクの作成
ストリームロードは HTTP 経由でデータをインポートします。このトピックでは、curl コマンドを使用してインポートジョブを送信する方法を示します。他の HTTP クライアントを使用することもできます。
構文
curl --location-trusted -u <username>:<password> -XPUT <url>
(
data_desc
)
[opt_properties]
-
Expectリクエストヘッダーでは、100-continue、つまり"Expect:100-continue"を指定することを推奨します。これにより、サーバーがインポートタスクのリクエストを拒否した場合に、不要なデータ転送を防ぎ、リソースのオーバーヘッドを削減できます。 -
StarRocks では、一部の単語は SQL 予約キーワードであり、SQL ステートメントで直接使用することはできません。SQL ステートメントで予約キーワードを使用するには、バックティック (`) で囲む必要があります。キーワードの詳細については、キーワードをご参照ください。
パラメーター
-
<username>:<password>:ご利用の StarRocks クラスターのユーザー名とパスワードです。このパラメーターは必須です。アカウントにパスワードがない場合は、<username>:を使用してください。 -
PUT:HTTP リクエストメソッドです。このメソッドは必須です。ストリームロードは PUT のみをサポートします。 -
<url>:宛先 StarRocks テーブルの URL です。このパラメーターは必須です。次の形式を使用してください:http://<fe_host>:<fe_http_port>/api/<database_name>/<table_name>/_stream_load。パラメーター
必須
説明
<fe_host>はい
StarRocks クラスター内の FE の IP アドレスです。
<fe_http_port>はい
FE の HTTP ポートです。デフォルト:18030。
ポート番号は、StarRocks の Services ページの Configure タブで http_port パラメーターを検索することで確認できます。また、
SHOW FRONTENDSコマンドを実行して、FE ノードの IP アドレスと HTTP ポートを表示することもできます。<database_name>はい
宛先テーブルを含むデータベースの名前です。
<table_name>はい
宛先 StarRocks テーブルの名前です。
-
desc:ソースデータファイルのプロパティ (ファイル名、フォーマット、列区切り文字、行区切り文字、ターゲットパーティション、StarRocks テーブルへの列マッピングなど) を記述します。形式は次のとおりです。-T <file_path> -H "format: CSV | JSON" -H "column_separator: <column_separator>" -H "row_delimiter: <row_delimiter>" -H "columns: <column1_name>[, <column2_name>, ... ]" -H "partitions: <partition1_name>[, <partition2_name>, ...]" -H "temporary_partitions: <temporary_partition1_name>[, <temporary_partition2_name>, ...]" -H "jsonpaths: [ \"<json_path1>\"[, \"<json_path2>\", ...] ]" -H "strip_outer_array: true | false" -H "json_root: <json_path>" -H "ignore_json_size: true | false" -H "compression: <compression_algorithm> | Content-Encoding: <compression_algorithm>"data_descのパラメーターは、共通パラメーター、CSV 用パラメーター、JSON 用パラメーターの 3 種類に分けられます。-
共通パラメーター
パラメーター
必須
説明
<file_path>はい
ソースデータファイルのパスです。ファイル名にはオプションで拡張子を含めることができます。
formatいいえ
データ形式です。有効な値:
CSVとJSON。デフォルト:CSV。partitionsいいえ
データをロードするパーティションです。このパラメーターが指定されていない場合、データはすべてのテーブルパーティションにロードされます。
temporary_partitionsいいえ
データをロードする一時パーティションです。
columnsいいえ
ソースデータファイルと StarRocks テーブル間の列マッピングを指定します。
-
ソースデータファイルの列が宛先テーブルの列の順序と一致する場合、このパラメーターは不要です。
-
ソースファイルの列がテーブルスキーマと一致しない場合は、このパラメーターを使用してマッピングを定義する必要があります。列は 2 つの方法でマッピングできます:
-
例 1:テーブルには
c1, c2, c3という列がありますが、ソースファイルの列の順序はc3, c2, c1です。-H "columns: c3, c2, c1"を指定する必要があります。 -
例 2:テーブルには
c1, c2, c3の 3 つの列があり、ソースファイルの最初の 3 つの列がそれらにマッピングされますが、ファイルには 4 番目の不要な列があります。この場合、-H "columns: c1, c2, c3, temp"と指定し、余分な列にはプレースホルダーを使用します。 -
例 3:テーブルには
year, month, dayの 3 つの列がありますが、ソースファイルには2018-06-01 01:02:03形式のタイムスタンプ列が 1 つしかありません。-H "columns: col, year = year(col), month=month(col), day=day(col)"を指定することで、テーブルの列を派生させることができます。
-
-
-
CSV 固有のパラメーター
パラメーター
必須
説明
column_separatorいいえ
ソースデータファイル内の列区切り文字です。デフォルト:
\t。印刷不能文字の場合は、
\xをプレフィックスとして付け、16 進値を使用します。たとえば、Hive の区切り文字\x01の場合は、-H "column_separator:\x01"と指定します。row_delimiterいいえ
ソースデータファイル内の行区切り文字です。デフォルト:
\n。重要curl を使用する場合、シェルが
\nを改行文字ではなくリテラル文字列として解釈する可能性があります。Bash では、
$'string'構文を使用して、\nや\tなどのエスケープシーケンスが正しく解釈されるようにします。例:-H $'row_delimiter:\n'。skip_headerいいえ
CSV ファイルの先頭でスキップするヘッダー行の数です。値は整数である必要があります。デフォルト:0。
このパラメーターを使用して、CSV ファイル内の列名や型などのメタデータ行を無視します。たとえば、このパラメーターを 1 に設定すると、最初の行がスキップされます。
ここで使用される行区切り文字は、ロードコマンドで指定されたものと一致する必要があります。
whereいいえ
ロードする行をフィルターするための条件です。
たとえば、k1 列が 20180601 に等しい行のみをロードするには、
-H "where: k1 = 20180601"と指定します。max_filter_ratioいいえ
データ品質エラーが原因で破棄できる行の最大比率です。デフォルト:
0(ゼロトレランス)。説明これは where 条件によってフィルターされた行には適用されません。
partitionsいいえ
ロードジョブの宛先パーティションです。
ターゲットパーティションを決定できる場合は、このパラメーターを指定することを推奨します。他のパーティションのデータは破棄されます。たとえば、パーティション p1 と p2 にデータをロードするには、
-H "partitions: p1, p2"と指定します。timeoutいいえ
ロードジョブのタイムアウト (秒単位) です。デフォルト:600。
値は 1 から 259200 の間でなければなりません。
strict_modeいいえ
このロードジョブの厳格モードを有効または無効にします。有効な値:
-
false(デフォルト):厳格モードは無効です。 -
true:厳格モードは有効です。有効にすると、ロード中の列の型変換に厳格なフィルタリングが適用されます。
timezoneいいえ
ロードジョブのタイムゾーンです。デフォルト:
Asia/Shanghai(UTC+08:00)。このパラメーターは、ロードで使用されるすべてのタイムゾーン関連関数に影響します。
exec_mem_limitいいえ
ロードジョブのメモリ制限です。デフォルト:2 GB。
-
-
JSON 固有のパラメーター
パラメーター
必須
説明
jsonpathsいいえ
JSON データから抽出するフィールドを指定します。これは、特定のソースフィールドをテーブル列にマッピングする必要がある場合にのみ必須です。値は JSON 形式である必要があります。
strip_outer_arrayいいえ
最も外側の配列構造を削除するかどうかを指定します。有効な値:
-
false(デフォルト):元の JSON 構造を保持します。JSON 配列全体が単一の値としてインポートされます。たとえば、データ
[{"k1" : 1, "k2" : 2},{"k1" : 3, "k2" : 4}]の場合、strip_outer_arrayをfalseに設定すると、データは単一の配列としてインポートされます。 -
true:ロードするデータが JSON 配列の場合、
strip_outer_arrayを true に設定する必要があります。たとえば、データ
[{"k1" : 1, "k2" : 2},{"k1" : 3, "k2" : 4}]の場合、strip_outer_arrayをtrueに設定すると、データは 2 つの別々の行としてインポートされます。
json_rootいいえ
インポートする JSON データのルート要素を指定します。このパラメーターは、フィールドを照合して JSON データをロードする場合にのみ必須です。値は有効な JsonPath 文字列である必要があります。デフォルト:空。これは JSON ファイル全体がインポートされることを意味します。
ignore_json_sizeいいえ
HTTP リクエスト内の JSON 本文のサイズチェックをスキップするかどうかを指定します。
説明デフォルトでは、HTTP リクエストの JSON 本文サイズは 100 MB を超えることはできません。サイズがこの制限を超えると、「
The size of this batch exceed the max size [104857600] of json type data data [8617627793]. Set ignore_json_size to skip check, although it may lead huge memory consuming.」というエラーが返されます。より大きなファイルをロードするには、リクエストにヘッダー"ignore_json_size:true"を追加して、このチェックをバイパスします。compression, Content-Encodingいいえ
リクエストボディに使用される圧縮アルゴリズムです。サポートされているアルゴリズム:GZIP、BZIP2、LZ4_FRAME、および ZSTD。
例:
curl --location-trusted -u root: -v 'http://127.0.0.1:18030/api/db0/tbl_simple/_stream_load' \-X PUT -H "expect:100-continue" \-H 'format: json' -H 'compression: lz4_frame' -T ./b.json.lz4。 -
-
-
opt_properties:インポートのオプションパラメーターを指定します。指定されたパラメーターは、インポートタスク全体に適用されます。形式は次のとおりです:
-H "label: <label_name>" -H "where: <condition1>[, <condition2>, ...]" -H "max_filter_ratio: <num>" -H "timeout: <num>" -H "strict_mode: true | false" -H "timezone: <string>" -H "load_mem_limit: <num>" -H "partial_update: true | false" -H "partial_update_mode: row | column" -H "merge_condition: <column_name>"パラメーター
必須
説明
labelいいえ
ロードジョブの一意のラベルです。StarRocks は、成功した複数のジョブで同じラベルを使用することを防ぎます。
ラベルを使用して、重複ロードを防ぐことができます。StarRocks は、成功したジョブのラベルを 30 分間保持します。
whereいいえ
フィルター条件を指定します。StarRocks は、
columnsパラメーターからの変換後にこのフィルターをデータに適用します。where句を満たすデータのみがインポートされます。たとえば、k1 列が 20180601 である行のみをインポートするには、
-H "where: k1 = 20180601"と指定します。max_filter_ratioいいえ
データ品質エラーが原因で破棄できる行の最大比率です。デフォルト:
0(ゼロトレランス)。説明これは
where条件によってフィルターされた行には適用されません。log_rejected_record_numいいえ
データ品質の問題により拒否された行のログに記録する最大数を指定します。v3.1 以降でサポートされています。有効な値:
0、-1、または正の整数。デフォルト:0。-
0:拒否された行はログに記録されません。 -
-1:拒否されたすべての行がログに記録されます。 -
正の整数、たとえば
n:各 BE (または CN) ノードで最大n個の拒否された行がログに記録されます。
timeoutいいえ
ロードジョブのタイムアウト (秒単位) です。デフォルト:600。
値は 1 から 259200 の間でなければなりません。
strict_modeいいえ
このロードジョブの厳格モードを有効にするかどうかを指定します。
-
false(デフォルト):厳格モードは無効です。 -
true:厳格モードは有効です。
timezoneいいえ
ロードジョブのタイムゾーンです。デフォルト:
Asia/Shanghai(UTC+08:00)。このパラメーターは、ロードで使用されるすべてのタイムゾーン関連関数に影響します。
load_mem_limitいいえ
ロードジョブのメモリ制限です。デフォルト:2 GB。
partial_updateいいえ
部分的な列の更新を有効にします。有効な値:
trueとfalse。デフォルト:false。partial_update_modeいいえ
部分更新のモードを指定します。有効な値:
rowとcolumn。-
row(デフォルト):行モード。少量のバッチで多くの列をリアルタイムに更新する場合に最適です。 -
column:列モード。多くの行にわたる少数の列の一括更新に最適です。このモードはパフォーマンスを大幅に向上させることができます。たとえば、テーブルのすべての行に対して 100 列中 10 列 (10%) を更新する場合、列モードはパフォーマンスを最大 10 倍向上させることができます。
merge_conditionいいえ
更新条件として使用する列名を指定します。行の更新は、この列の新しいデータの値が既存の値以上である場合にのみ発生します。
説明指定された列はプライマリキーの一部であってはなりません。この機能はプライマリキーテーブルでのみ使用可能です。
-
例
この例では、StarRocks クラスターの load_test データベースにある example_table テーブルに data.csv をロードします。完全な例については、「完全なデータロードの例」をご参照ください。
curl --location-trusted -u "root:" \
-H "Expect:100-continue" \
-H "label:label2" \
-H "column_separator: ," \
-T data.csv -XPUT \
http://172.17.**.**:18030/api/load_test/example_table/_stream_load
戻り値
ロードタスクが完了すると、結果が JSON 形式で返されます。以下はサンプルです。
{
"TxnId": 9,
"Label": "label2",
"Status": "Success",
"Message": "OK",
"NumberTotalRows": 4,
"NumberLoadedRows": 4,
"NumberFilteredRows": 0,
"NumberUnselectedRows": 0,
"LoadBytes": 45,
"LoadTimeMs": 235,
"BeginTxnTimeMs": 101,
"StreamLoadPlanTimeMs": 102,
"ReadDataTimeMs": 0,
"WriteDataTimeMs": 11,
"CommitAndPublishTimeMs": 19
}
戻り値には次のパラメーターが含まれます:
|
パラメーター |
説明 |
|
|
ロードタスクのトランザクション ID です。 |
|
|
ロードタスクのラベルです。 |
|
|
ロードステータスです。有効な値は次のとおりです:
|
|
|
このフィールドは、すでに存在するラベルに対応するインポートタスクのステータスを示します。このフィールドは、
|
|
|
ロードステータスの詳細です。ロードが失敗した場合、このパラメーターは理由を提供します。 |
|
|
データストリームから読み取られた総行数です。 |
|
|
ロードされた行数です。このパラメーターは、 |
|
|
データ品質が低いために除外された行数です。 |
|
|
|
|
|
ソースファイルのサイズ (バイト単位) です。 |
|
|
ロードタスクの持続時間 (ミリ秒単位) です。 |
|
|
トランザクションの開始にかかった時間 (ミリ秒単位) です。 |
|
|
実行計画の生成にかかった時間 (ミリ秒単位) です。 |
|
|
データの読み取りにかかった時間 (ミリ秒単位) です。 |
|
|
データの書き込みにかかった時間 (ミリ秒単位) です。 |
|
|
ロードが失敗した場合にのみ返されます。
たとえば、エラー行情報をダウンロードするには:
エクスポートされたエラー行は、 エラー情報に基づいてロードタスクを調整し、再試行してください。 |
ロードタスクのキャンセル
ストリームロードタスクは手動でキャンセルできません。タスクがタイムアウトした場合やインポートエラーが発生した場合、システムは自動的にタスクをキャンセルします。応答の ErrorURL を使用して、トラブルシューティングのためにエラーの詳細をダウンロードしてください。
データロードの例
この例では、curl コマンドを使用してロードタスクを送信する方法を示します。
-
宛先テーブルの作成
-
SSH を使用して StarRocks クラスターのマスターノードにログインします。詳細については、「クラスターへのログイン」をご参照ください。
-
次のコマンドを実行して、MySQL クライアントを使用して StarRocks クラスターに接続します。
mysql -h127.0.0.1 -P 9030 -uroot -
次のコマンドを実行して、データベースとテーブルを作成します。
CREATE DATABASE IF NOT EXISTS load_test; USE load_test; CREATE TABLE IF NOT EXISTS example_table ( id INT, name VARCHAR(50), age INT ) DUPLICATE KEY(id) DISTRIBUTED BY HASH(id) BUCKETS 3 PROPERTIES ( "replication_num" = "1" -- レプリケーション数を 1 に設定します。 );Ctrl+Dを押して MySQL クライアントを終了します。
-
-
テストデータの準備
CSV データ
たとえば、次の内容で
data.csvという名前のファイルを作成します:id,name,age 1,Alice,25 2,Bob,30 3,Charlie,35JSON データ
たとえば、次の内容で
json.dataという名前のファイルを作成します:{"id":1,"name":"Emily","age":25} {"id":2,"name":"Benjamin","age":35} {"id":3,"name":"Olivia","age":28} {"id":4,"name":"Alexander","age":60} {"id":5,"name":"Ava","age":17} -
次のコマンドを実行してロードタスクを送信します。
CSV データ
curl --location-trusted -u "root:" \ -H "Expect:100-continue" \ -H "label:label1" \ -H "column_separator: ," \ -T data.csv -XPUT \ http://172.17.**.**:18030/api/load_test/example_table/_stream_loadJSON データ
curl --location-trusted -u "root:" \ -H "Expect:100-continue" \ -H "label:label2" \ -H "format:json" \ -T json.data -XPUT \ http://172.17.**.**:18030/api/load_test/example_table/_stream_load
コード統合の例
-
Java でのストリームロードジョブの開発。詳細については、stream_load をご参照ください。
-
Spark とのストリームロードの統合。詳細については、01_sparkStreaming2StarRocks をご参照ください。