Spark コネクタを使用したデータのロード (推奨)
StarRocks は、Apache Spark™ 用の StarRocks Connector(以下、Spark コネクタ)という独自開発のコネクタを提供しており、Spark を使用して StarRocks テーブルにデータをロードするのに役立ちます。基本的な原則は、データを蓄積し、STREAM LOAD を通じて一度にすべてのデータを StarRocks にロードすることです。Spark コネクタは Spark DataSource V2 に基づいて実装されています。DataSource は Spark DataFrames または Spark SQL を使用して作成できます。バッチモードと構造化ストリーミングモードの両方がサポートされています。
注意
StarRocks テーブルに対して SELECT および INSERT 権限を持つユーザーのみが、このテーブルにデータをロードできます。GRANT の指示に従って、これらの権限をユーザーに付与できます。
バージョン要件
| Spark コネクタ | Spark | StarRocks | Java | Scala |
|---|---|---|---|---|
| 1.1.2 | 3.2, 3.3, 3.4, 3.5 | 2.5 以降 | 8 | 2.12 |
| 1.1.1 | 3.2, 3.3, または 3.4 | 2.5 以降 | 8 | 2.12 |
| 1.1.0 | 3.2, 3.3, または 3.4 | 2.5 以降 | 8 | 2.12 |
注意
- Spark コネクタのバージョン間の動作の変更については、Upgrade Spark connector を参照してください。
- Spark コネクタはバージョン 1.1.1 以降、MySQL JDBC ドライバを提供していないため、ドライバを手動で spark クラスパスにインポートする必要があります。ドライバは MySQL サイト または Maven Central で見つけることができます。
Spark コネクタの取得
Spark コネクタ JAR ファイルを取得する方法は以下の通りです:
- コンパイル済みの Spark Connector JAR ファイルを直接ダウンロードします。
- Maven プロジェクトに Spark コネクタを依存関係として追加し、JAR ファイルをダウンロードします。
- Spark Connector のソースコードを自分でコンパイルして JAR ファイルを作成します。
Spark コネクタ JAR ファイルの命名形式は starrocks-spark-connector-${spark_version}_${scala_version}-${connector_version}.jar です。
例えば、Spark 3.2 と Scala 2.12 を環境にインストールし、Spark コネクタ 1.1.0 を使用したい場合、starrocks-spark-connector-3.2_2.12-1.1.0.jar を使用できます。
注意
一般に、最新バージョンの Spark コネクタは Spark の直近3つのバージョンとの互換性のみを維持しています。
コンパイル済み Jar ファイルのダウンロード
Maven Central Repository から対応するバージョンの Spark コネクタ JAR を直接ダウンロードします。
Maven 依存関係
-
Maven プロジェクトの
pom.xmlファイルに、以下の形式で Spark コネクタを依存関係として追加します。spark_version、scala_version、connector_versionをそれぞれのバージョンに置き換えてください。<dependency>
<groupId>com.starrocks</groupId>
<artifactId>starrocks-spark-connector-${spark_version}_${scala_version}</artifactId>
<version>${connector_version}</version>
</dependency> -
例えば、環境の Spark バージョンが 3.2、Scala バージョンが 2.12 で、Spark コネクタ 1.1.0 を選択する場合、以下の依存関係を追加する必要があります:
<dependency>
<groupId>com.starrocks</groupId>
<artifactId>starrocks-spark-connector-3.2_2.12</artifactId>
<version>1.1.0</version>
</dependency>
自分でコンパイル
-
Spark コネクタパッケージ をダウンロードします。
-
以下のコマンドを実行して、Spark コネクタのソースコードを JAR ファイルにコンパイルします。
spark_versionは対応する Spark バージョンに置き換えてください。sh build.sh <spark_version>例えば、環境の Spark バージョンが 3.2 の場合、以下のコマンドを実行する必要があります:
sh build.sh 3.2 -
コンパイル後に生成された Spark コネクタ JAR ファイルを
target/ディレクトリで見つけます。例えば、starrocks-spark-connector-3.2_2.12-1.1.0-SNAPSHOT.jarです。
注意
正式にリリースされていない Spark コネクタの名前には
SNAPSHOTサフィックスが含まれます。
パラメータ
| パラメータ | 必須 | デフォルト値 | 説明 |
|---|---|---|---|
| starrocks.fe.http.url | YES | None | StarRocks クラスター内の FE の HTTP URL。複数の URL を指定でき、カンマ (,) で区切る必要が あります。形式: <fe_host1>:<fe_http_port1>,<fe_host2>:<fe_http_port2>。バージョン 1.1.1 以降、URL に http:// プレフィックスを追加することもできます。例: http://<fe_host1>:<fe_http_port1>,http://<fe_host2>:<fe_http_port2>。 |
| starrocks.fe.jdbc.url | YES | None | FE の MySQL サーバーに接続するために使用されるアドレス。形式: jdbc:mysql://<fe_host>:<fe_query_port>。 |
| starrocks.table.identifier | YES | None | StarRocks テーブルの名前。形式: <database_name>.<table_name>。 |
| starrocks.user | YES | None | StarRocks クラスターアカウントのユーザー名。ユーザーは StarRocks テーブルに対して SELECT および INSERT 権限 を持っている必要があります。 |
| starrocks.password | YES | None | StarRocks クラスターアカウントのパスワード。 |
| starrocks.write.label.prefix | NO | spark- | Stream Load で使用されるラベルプレフィックス。 |
| starrocks.write.enable.transaction-stream-load | NO | TRUE | データをロードするために Stream Load トランザクションインターフェース を使用するかどうか。StarRocks v2.5 以降が必要です。この機能は、トランザクションでより多くのデータを少ないメモリ使用量でロードし、パフォーマンスを向上させることができます。 注意: バージョン 1.1.1 以降、このパラメータは starrocks.write.max.retries の値が非正の場合に のみ有効です。なぜなら、Stream Load トランザクションインターフェースはリトライをサポートしていないからです。 |
| starrocks.write.buffer.size | NO | 104857600 | StarRocks に一度に送信される前にメモリに蓄積できるデータの最大サイズ。このパラメータを大きな値に設定すると、ロードパフォーマンスが向上しますが、ロード遅延が増加する可能性があります。 |
| starrocks.write.buffer.rows | NO | Integer.MAX_VALUE | バージョン 1.1.1 以降でサポートされています。StarRocks に一度に送信される前にメモリに蓄積できる行の最大数。 |
| starrocks.write.flush.interval.ms | NO | 300000 | StarRocks にデータを送信する間隔。このパラメータはロード遅延を制御するために使用されます。 |
| starrocks.write.max.retries | NO | 3 | バージョン 1.1.1 以降でサポートされています。ロードが失敗した場合に同じデータバッチに対して Stream Load を実行するためにコネクタがリトライする回数。 注意: Stream Load トランザクションインターフェースはリトライをサポートしていないため、このパラメータが正の場合、コネクタは常に Stream Load インターフェースを使用し、 starrocks.write.enable.transaction-stream-load の値を無視します。 |
| starrocks.write.retry.interval.ms | NO | 10000 | バージョン 1.1.1 以降でサポートされています。ロードが失敗した場合に同じデータバッチに対して Stream Load をリトライする間隔。 |
| starrocks.columns | NO | None | データをロードしたい StarRocks テーブルの列。複数の列を指定でき、カンマ (,) で区切る必要があります。例: "col0,col1,col2"。 |
| starrocks.column.types | NO | None | バージョン 1.1.1 以降でサポートされています。StarRocks テーブルと デフォルトマッピング から推測されるデフォルトを使用する代わりに、Spark 用の列データ型をカスタマイズします。パラメータ値は Spark の StructType#toDDL の出力と同じ DDL 形式のスキーマです。例: col0 INT, col1 STRING, col2 BIGINT。カスタマイズが必要な列のみを指定する必要があります。使用例としては、BITMAP または HLL 型の列にデータをロードすることです。 |
| starrocks.write.properties.* | NO | None | Stream Load の動作を制御するために使用されるパラメータ。例えば、パラメータ starrocks.write.properties.format はロードするデータの形式を指定します。例: CSV または JSON。サポートされているパラメータとその説明のリストについては、STREAM LOAD を参照してください。 |
| starrocks.write.properties.format | NO | CSV | Spark コネクタが StarRocks にデータを送信する前に各データバッチを変換する際のファイル形式。有効な値: CSV および JSON。 |
| starrocks.write.properties.row_delimiter | NO | \n | CSV 形式のデータの行区切り文字。 |
| starrocks.write.properties.column_separator | NO | \t | CSV 形式のデータの列区切り文字。 |
| starrocks.write.num.partitions | NO | None | Spark がデータを書き込む際に並列で使用するパーティション数。データ量が少ない場合、パーティション数を減らしてロードの同時実行性と頻度を下げることができます。このパラメータのデフォルト値は Spark によって決定されます。ただし、この方法は Spark Shuffle コストを引き起こす可能性があります。 |
| starrocks.write.partition.columns | NO | None | Spark のパーティション列。このパラメータは starrocks.write.num.partitions が指定されている場合にのみ有効です。このパラメータが指定されていない場合、書き込まれるすべての列がパーティション分割に使用されます。 |
| starrocks.timezone | NO | JVM のデフォルトタイムゾーン | バージョン 1.1.1 以降でサポートされています。Spark の TimestampType を StarRocks の DATETIME に変換する際に使用されるタイムゾーン。デフォルトは ZoneId#systemDefault() によって返される JVM のタイムゾーンです。形式は Asia/Shanghai のようなタイムゾーン名、または +08:00 のようなゾーンオフセットです。 |
Spark と StarRocks 間のデータ型マッピング
-
デ フォルトのデータ型マッピングは次の通りです:
Spark データ型 StarRocks データ型 BooleanType BOOLEAN ByteType TINYINT ShortType SMALLINT IntegerType INT LongType BIGINT StringType LARGEINT FloatType FLOAT DoubleType DOUBLE DecimalType DECIMAL StringType CHAR StringType VARCHAR StringType STRING StringType JSON DateType DATE TimestampType DATETIME ArrayType ARRAY
注意:
バージョン 1.1.1 以降でサポートされています。詳細な手順については、Load data into columns of ARRAY type を参照してください。 -
データ型マッピングをカスタマイズすることもできます。
例えば、StarRocks テーブルに BITMAP および HLL 列が含まれている場合、Spark はこれらのデータ型をサポートしていません。Spark で対応するデータ型をカスタマイズする必要があります。詳細な手順については、BITMAP および HLL 列にデータをロードする方法を参照してください。BITMAP および HLL はバージョン 1.1.1 以降でサポートされています。