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

E-MapReduce:Stream load

最終更新日:Aug 22, 2026

ストリームロードを使用して、ローカルファイルまたはデータストリームから 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
}

戻り値には次のパラメーターが含まれます:

パラメーター

説明

TxnId

ロードタスクのトランザクション ID です。

Label

ロードタスクのラベルです。

Status

ロードステータスです。有効な値は次のとおりです:

  • Success:データは正常にロードされ、表示可能です。

  • Publish Timeout:ロードタスクはコミットされましたが、データの可視性が遅れる可能性があります。タスクをリトライしないでください。

  • Label Already Exists:Label がすでに存在することを示します。Label を変更する必要があります。

  • Fail:ロードは失敗しました。

ExistingJobStatus

このフィールドは、すでに存在するラベルに対応するインポートタスクのステータスを示します。このフィールドは、Status が Label Already Exists の場合にのみ表示されます。このステータスを確認して、インポートタスクの状態を判断できます。

  • RUNNING:タスクは実行中です。

  • FINISHED:タスクは完了しました。

Message

ロードステータスの詳細です。ロードが失敗した場合、このパラメーターは理由を提供します。

NumberTotalRows

データストリームから読み取られた総行数です。

NumberLoadedRows

ロードされた行数です。このパラメーターは、Status が Success の場合にのみ適用されます。

NumberFilteredRows

データ品質が低いために除外された行数です。

NumberUnselectedRows

where 条件によって除外された行数です。

LoadBytes

ソースファイルのサイズ (バイト単位) です。

LoadTimeMs

ロードタスクの持続時間 (ミリ秒単位) です。

BeginTxnTimeMs

トランザクションの開始にかかった時間 (ミリ秒単位) です。

StreamLoadPlanTimeMs

実行計画の生成にかかった時間 (ミリ秒単位) です。

ReadDataTimeMs

データの読み取りにかかった時間 (ミリ秒単位) です。

WriteDataTimeMs

データの書き込みにかかった時間 (ミリ秒単位) です。

ErrorURL

ロードが失敗した場合にのみ返されます。

ErrorURL を使用して、データ品質チェックに失敗したためにインポートプロセス中に除外されたエラー行の詳細を表示できます。インポートタスクを送信する際に、オプションのパラメーター log_rejected_record_num を使用して、ログに記録するエラー行の最大数を指定できます。

curl "url" コマンドを実行してエラー行の情報を直接表示するか、wget "url" コマンドを実行して情報をエクスポートできます。

たとえば、エラー行情報をダウンロードするには:

wget "http://172.17.**.**:18040/api/_load_error_log?file=error_log_b74dccdcf0ceb4de_e82b2709c6c013ad"

エクスポートされたエラー行は、_load_error_log?file=error_log_b74dccdcf0ceb4de_e82b2709c6c013ad という名前のローカルファイルに保存されます。cat _load_error_log?file=error_log_b74dccdcf0ceb4de_e82b2709c6c013ad コマンドを使用して、ファイルの内容を表示できます。

エラー情報に基づいてロードタスクを調整し、再試行してください。

ロードタスクのキャンセル

ストリームロードタスクは手動でキャンセルできません。タスクがタイムアウトした場合やインポートエラーが発生した場合、システムは自動的にタスクをキャンセルします。応答の ErrorURL を使用して、トラブルシューティングのためにエラーの詳細をダウンロードしてください。

データロードの例

この例では、curl コマンドを使用してロードタスクを送信する方法を示します。

  1. 宛先テーブルの作成

    1. SSH を使用して StarRocks クラスターのマスターノードにログインします。詳細については、「クラスターへのログイン」をご参照ください。

    2. 次のコマンドを実行して、MySQL クライアントを使用して StarRocks クラスターに接続します。

      mysql -h127.0.0.1  -P 9030 -uroot
    3. 次のコマンドを実行して、データベースとテーブルを作成します。

      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 クライアントを終了します。

  2. テストデータの準備

    CSV データ

    たとえば、次の内容で data.csv という名前のファイルを作成します:

    id,name,age
    1,Alice,25
    2,Bob,30
    3,Charlie,35

    JSON データ

    たとえば、次の内容で 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}
  3. 次のコマンドを実行してロードタスクを送信します。

    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_load

    JSON データ

    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 をご参照ください。