Sparkは、大規模データ処理のための統合分析エンジンです。Hologresは、コミュニティ版 Spark および EMR Serverless Spark と効率的に統合することで、データウェアハウスを迅速に構築できます。Hologres の Spark コネクタは、Spark クラスターにおける Hologres カタログの作成に対応しています。これにより、外部テーブルを使用した高性能なバッチ読み込みとインポートが可能になり、ネイティブ JDBC よりも優れたパフォーマンスを実現します。
制限事項
Spark コネクタは、Hologres バージョン 1.3 以降が必要です。インスタンスのバージョンは、Hologres コンソールの インスタンスの詳細 ページで確認できます。インスタンスのバージョンが 1.3 より古い場合は、インスタンスをアップグレードするか、Hologres DingTalk グループに参加してアップグレードをリクエストしてください。詳細については、「オンラインサポートの追加利用方法」をご参照ください。
前提条件
-
spark-sql、spark-shell、またはpysparkコマンドを実行できる Spark 環境を準備してください。依存関係の問題を回避し、より多くの機能を利用するには、Spark 3.3.0 以降を使用してください。-
Alibaba Cloud EMR Spark を使用して、Spark 環境を迅速にセットアップし、Hologres インスタンスに接続できます。詳細については、「EMR Spark の機能」をご参照ください。
-
または、独立した Spark 環境をセットアップすることもできます。詳細については、「Apache Spark」をご参照ください。
-
-
Spark で Hologres の読み取りと書き込みを行うには、
hologres-connector-spark-3.xコネクターが必要です。本トピックでは、Maven 中央リポジトリからダウンロードできるバージョン 1.5.2 を例として使用します。このコネクターはオープンソースです。詳細については、「Hologres-Connectors」をご参照ください。 -
IntelliJ IDEA などの IDE を使用して Java で Spark ジョブを開発し、ローカルでデバッグするには、pom.xml ファイルに以下の Maven 依存関係を追加してください。
<dependency> <groupId>com.alibaba.hologres</groupId> <artifactId>hologres-connector-spark-3.x</artifactId> <version>1.5.2</version> <classifier>jar-with-dependencies</classifier> </dependency>
Hologres カタログ
Hologres コネクタ 1.5.2 以降は Hologres カタログをサポートしており、外部テーブルを使用して Hologres への読み書きができます。
Spark の各 Hologres カタログは、Hologres のデータベースに対応します。Hologres カタログの各名前空間は、対応するデータベースのスキーマに対応します。以降のセクションでは、Spark で Hologres カタログを使用する方法について説明します。
Hologres カタログはテーブルの作成をサポートしていません。
このトピックでは、Hologres インスタンス内の以下のデータベースとテーブルを使用します:
test_db -- データベース
public.test_table1 -- public スキーマ内のテーブル
public.test_table2
test_schema.test_table3 -- test_schema スキーマ内のテーブル
Hologres カタログの初期化
Spark クラスターで spark-sql を起動し、Hologres コネクタをロードして、カタログパラメータを指定します。
spark-sql --jars hologres-connector-spark-3.x-1.5.2-jar-with-dependencies.jar \
--conf spark.sql.catalog.hologres_external_test_db=com.alibaba.hologres.spark3.HoloTableCatalog \
--conf spark.sql.catalog.hologres_external_test_db.username=*** \
--conf spark.sql.catalog.hologres_external_test_db.password=*** \
--conf spark.sql.catalog.hologres_external_test_db.jdbcurl=jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db
Hologres カタログコマンド
-
Hologres カタログのロード
Spark の Hologres カタログは、Hologres のデータベースに対応します。このマッピングは、セッションの間は固定されます。
USE hologres_external_test_db; -
すべての名前空間のクエリ
Spark の名前空間は Hologres のスキーマに対応します。デフォルトのスキーマは public です。デフォルトのスキーマを変更するには、
USEコマンドを使用します。-- Hologres カタログ内のすべての名前空間 (Hologres データベースのスキーマに対応) を表示します。 SHOW NAMESPACES; -
名前空間内のテーブルのクエリ
-
すべてのテーブルのクエリ
SHOW TABLES; -
特定の名前空間内のテーブルのクエリ
USE test_schema; SHOW TABLES; -- または、次のステートメントも使用できます。 SHOW TABLES IN test_schema;
-
-
テーブルの読み書き
SELECT および INSERT ステートメントを使用して、外部テーブルへの読み書きを行います。
-- テーブルからデータを読み取ります。 SELECT * FROM public.test_table1; -- テーブルにデータを書き込みます。 INSERT INTO test_schema.test_table3 SELECT * FROM public.test_table1;
Hologres へのデータのインポート
このセクションのテストデータは、TPC-H データセットの customer テーブルに基づいています。Spark は CSV ファイルからデータを読み取り、Hologres テーブルにデータを書き込むことができます。サンプルの customer データをダウンロードできます。次の SQL ステートメントで、customer_holo_table を作成します。
CREATE TABLE customer_holo_table
(
c_custkey BIGINT ,
c_name TEXT ,
c_address TEXT ,
c_nationkey INT ,
c_phone TEXT ,
c_acctbal DECIMAL(15,2) ,
c_mktsegment TEXT ,
c_comment TEXT
);
Spark-SQL を使用したインポート
Spark-SQL では、カタログを使用して Hologres テーブルのメタデータをロードするほうが便利です。また、一時テーブルを作成して Hologres テーブルを宣言することもできます。
-
1.5.2 より前のバージョンの Hologres Spark コネクタはカタログをサポートしていません。一時テーブルを作成して Hologres テーブルを宣言することしかできません。
-
Hologres Spark コネクタのパラメーターの詳細については、「parameters」をご参照ください。
-
Hologres カタログを初期化します。
spark-sql --jars hologres-connector-spark-3.x-1.5.2-jar-with-dependencies.jar \ --conf spark.sql.catalog.hologres_external_test_db=com.alibaba.hologres.spark3.HoloTableCatalog \ --conf spark.sql.catalog.hologres_external_test_db.username=*** \ --conf spark.sql.catalog.hologres_external_test_db.password=*** \ --conf spark.sql.catalog.hologres_external_test_db.jdbcurl=jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db -
CSV データソースから Hologres テーブルにデータをインポートします。
説明Spark の INSERT INTO 構文では、
column_listを使用して列のサブセットを指定することはできません。たとえば、INSERT INTO hologresTable(c_custkey) SELECT c_custkey FROM csvTableを使用して、c_custkey フィールドのみにデータを書き込むことはできません。特定のフィールドにデータを書き込みたい場合は、
CREATE TEMPORARY VIEWステートメントを使用して、必要なフィールドのみを含む Hologres の一時ビューを宣言します。カタログを使用
-- Hologres カタログをロードします。 USE hologres_external_test_db; -- CSV データソースを作成します。 CREATE TEMPORARY VIEW csvTable ( c_custkey BIGINT, c_name STRING, c_address STRING, c_nationkey INT, c_phone STRING, c_acctbal DECIMAL(15, 2), c_mktsegment STRING, c_comment STRING) USING csv OPTIONS ( path "resources/customer", sep "," -- ローカルテストでは、ファイルの絶対パスを使用します。 ); -- CSV テーブルのデータを Hologres に書き込みます。 INSERT INTO public.customer_holo_table SELECT * FROM csvTable;一時ビューを使用
-- CSV データソースを作成します。 CREATE TEMPORARY VIEW csvTable ( c_custkey BIGINT, c_name STRING, c_address STRING, c_nationkey INT, c_phone STRING, c_acctbal DECIMAL(15, 2), c_mktsegment STRING, c_comment STRING) USING csv OPTIONS ( path "resources/customer", sep "," ); -- Hologres の一時ビューを作成します。 CREATE TEMPORARY VIEW hologresTable ( c_custkey BIGINT, c_name STRING, c_phone STRING) USING hologres OPTIONS ( jdbcurl "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db", username "***", password "***", table "customer_holo_table" ); INSERT INTO hologresTable SELECT c_custkey,c_name,c_phone FROM csvTable;
DataFrame を使用したインポート
spark-shell や pyspark などのツールを使用して Spark ジョブを開発し、write API を呼び出してデータを書き込むことができます。このジョブは CSV ファイルからデータを読み取り、DataFrame に変換した後、DataFrame を Hologres インスタンスに書き込みます。以降のセクションでは、プログラミング言語ごとのサンプルコードを示します。Hologres Spark コネクタのパラメーターの詳細については、「parameters」をご参照ください。
Scala
import org.apache.spark.sql.types._
import org.apache.spark.sql.SaveMode
// CSV データソースのスキーマ。
val schema = StructType(Array(
StructField("c_custkey", LongType),
StructField("c_name", StringType),
StructField("c_address", StringType),
StructField("c_nationkey", IntegerType),
StructField("c_phone", StringType),
StructField("c_acctbal", DecimalType(15, 2)),
StructField("c_mktsegment", StringType),
StructField("c_comment", StringType)
))
// CSV ファイルからデータを読み取り、DataFrame に格納します。
val csvDf = spark.read.format("csv").schema(schema).option("sep", ",").load("resources/customer")
// DataFrame を Hologres に書き込みます。
csvDf.write
.format("hologres")
.option("username", "***")
.option("password", "***")
.option("jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db")
.option("table", "customer_holo_table")
.mode(SaveMode.Append)
.save()
Java
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.types.*;
import org.apache.spark.sql.SaveMode;
import java.util.Arrays;
import java.util.List;
public class SparkTest {
public static void main(String[] args) {
// CSV データソースのスキーマ。
List<StructField> asList =
Arrays.asList(
DataTypes.createStructField("c_custkey", DataTypes.LongType, true),
DataTypes.createStructField("c_name", DataTypes.StringType, true),
DataTypes.createStructField("c_address", DataTypes.StringType, true),
DataTypes.createStructField("c_nationkey", DataTypes.IntegerType, true),
DataTypes.createStructField("c_phone", DataTypes.StringType, true),
DataTypes.createStructField("c_acctbal", new DecimalType(15, 2), true),
DataTypes.createStructField("c_mktsegment", DataTypes.StringType, true),
DataTypes.createStructField("c_comment", DataTypes.StringType, true));
StructType schema = DataTypes.createStructType(asList);
// ローカルモードで実行します。
SparkSession spark = SparkSession.builder()
.appName("Spark CSV Example")
.master("local[*]")
.getOrCreate();
// CSV ファイルからデータを読み取り、DataFrame に格納します。
// ローカルテストでは、customer データの絶対パスを使用します。
Dataset<Row> csvDf = spark.read().format("csv").schema(schema).option("sep", ",").load("resources/customer");
// DataFrame を Hologres に書き込みます。
csvDf.write.format("hologres").option(
"username", "***").option(
"password", "***").option(
"jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db").option(
"table", "customer_holo_table").mode(
"append").save();
}
}
pom.xml ファイルに次の依存関係を追加します。
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.13</artifactId>
<version>3.5.4</version>
<scope>provided</scope>
</dependency>
Python
from pyspark.sql.types import *
# CSV データソースのスキーマ。
schema = StructType([
StructField("c_custkey", LongType()),
StructField("c_name", StringType()),
StructField("c_address", StringType()),
StructField("c_nationkey", IntegerType()),
StructField("c_phone", StringType()),
StructField("c_acctbal", DecimalType(15, 2)),
StructField("c_mktsegment", StringType()),
StructField("c_comment", StringType())
])
# CSV ファイルからデータを読み取り、DataFrame に格納します。
csvDf = spark.read.csv("resources/customer", header=False, schema=schema, sep=',')
# DataFrame を Hologres に書き込みます。
csvDf.write.format("hologres").option(
"username", "***").option(
"password", "***").option(
"jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db").option(
"table", "customer_holo_table").mode(
"append").save()
各言語で Spark ジョブを実行するには、次の手順に従ってください:
-
Scala
-
サンプルコードを使用して
sparktest.scalaファイルを作成し、次のコマンドを実行してジョブを実行します。-- 依存関係をロードします。 spark-shell --jars hologres-connector-spark-3.x-1.5.2-jar-with-dependencies.jar -- ローカルテストでは、ファイルをロードするために絶対パスを使用します。 scala> :load D:/sparktest.scala -
または、依存関係をロードした後に、サンプルコードをシェルに直接貼り付けて実行することもできます。
-
-
Java
開発ツールを使用してサンプルコードをインポートし、Maven でパッケージ化します。たとえば、出力される JAR ファイルが
spark_test.jarの場合は、次のコマンドを実行してジョブを実行します。-- ジョブの JAR ファイルを絶対パスで指定します。 spark-submit --class SparkTest --jars hologres-connector-spark-3.x-1.5.2-jar-with-dependencies.jar D:\spark_test.jar -
Python
次のコマンドを実行した後、サンプルコードをシェルに貼り付けて実行できます。
pyspark --jars hologres-connector-spark-3.x-1.5.2-jar-with-dependencies.jar
Hologresからのデータ読み取り
-
バージョン 1.3.2 以降、Spark コネクタは Hologres からのデータ読み取りに対応しています。デフォルトの Spark
jdbc-connectorと比較して、spark-connectorは Hologres テーブルのシャードに基づいてデータを並列で読み取るため、より優れたパフォーマンスを発揮します。読み取りの並列度は、テーブル内のシャード数に関連します。spark-connectorは、read.max_task_countパラメーターで並列度を制限できます。最終的にジョブはMin(shardCount, max_task_count)個の読み取りタスクを生成します。また、スキーマ推論にも対応しており、スキーマを指定しない場合は、Hologres テーブルスキーマから Spark スキーマを推論します。 -
Spark コネクタのバージョン 1.5.0 以降、Hologres テーブルからのデータ読み取りは、述語プッシュダウン、LIMIT プッシュダウン、およびカラムプルーニングに対応しています。Hologres の
SELECT QUERYを使用してデータを読み取ることもできます。このバージョンではバッチ読み取りモードが導入されており、以前のバージョンと比較して読み取りパフォーマンスが 3~4 倍向上します。
Spark SQLによるデータ読み取り
Spark SQL を使用する場合、カタログで Hologres テーブルのメタデータをロードするか、一時テーブルを作成して Hologres テーブルを宣言できます。
-
Hologres Spark コネクタのバージョン 1.5.2 より前は、カタログに対応していません。Hologres テーブルは、一時テーブルを作成することによってのみ宣言できます。
-
Hologres Spark コネクタのパラメーターの詳細については、「パラメーター」をご参照ください。
-
Hologres カタログを初期化します。
spark-sql --jars hologres-connector-spark-3.x-1.5.2-jar-with-dependencies.jar \ --conf spark.sql.catalog.hologres_external_test_db=com.alibaba.hologres.spark3.HoloTableCatalog \ --conf spark.sql.catalog.hologres_external_test_db.username=*** \ --conf spark.sql.catalog.hologres_external_test_db.password=*** \ --conf spark.sql.catalog.hologres_external_test_db.jdbcurl=jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db -
Hologres からデータを読み取ります。
-
カタログを使用してデータを読み取ります。
-- Hologres カタログをロードします。 USE hologres_external_test_db; -- Hologres テーブルからデータを読み取ります。カラムプルーニングと述語プッシュダウンに対応しています。 SELECT c_custkey,c_name,c_phone FROM public.customer_holo_table WHERE c_custkey < 500 LIMIT 10; -
一時テーブルを作成してデータを読み取ります。
テーブル
CREATE TEMPORARY VIEW hologresTable USING hologres OPTIONS ( jdbcurl "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db", username "***", password "***", read.max_task_count "80", -- Hologres テーブルから読み取るタスクの最大数。 table "customer_holo_table" ); -- カラムプルーニングと述語プッシュダウンに対応しています。 SELECT c_custkey,c_name,c_phone FROM hologresTable WHERE c_custkey < 500 LIMIT 10;クエリ
CREATE TEMPORARY VIEW hologresTable USING hologres OPTIONS ( jdbcurl "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db", username "***", password "***", read.query "SELECT c_custkey,c_name,c_phone FROM customer_holo_table WHERE c_custkey < 500 LIMIT 10" ); SELECT * FROM hologresTable LIMIT 5;
-
Hologresデータの DataFrameへの読み込み
spark-shell や pyspark などのツールで Spark ジョブを開発する際は、Spark の read API を呼び出してデータを DataFrame にロードできます。以下の例では、さまざまなプログラミング言語で Hologres テーブルから DataFrame にデータを読み取る方法を説明します。Hologres Spark コネクタのパラメーターの詳細については、「パラメーター」をご参照ください。
Scala
val readDf = (
spark.read
.format("hologres")
.option("username", "***")
.option("password", "***")
.option("jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db")
.option("table", "customer_holo_table")
.option("read.max_task_count", "80") // Hologres テーブルから読み取るタスクの最大数。
.load()
.filter("c_custkey < 500")
)
readDf.select("c_custkey", "c_name", "c_phone").show(10)
Java
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
public class SparkSelect {
public static void main(String[] args) {
// ローカルモードで実行します。
SparkSession spark = SparkSession.builder()
.appName("Spark Hologres Example")
.master("local[*]")
.getOrCreate();
Dataset<Row> readDf = (
spark.read
.format("hologres")
.option("username", "***")
.option("password", "***")
.option("jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db")
.option("table", "customer_holo_table")
.option("read.max_task_count", "80") // Hologres テーブルから読み取るタスクの最大数。
.load()
.filter("c_custkey < 500")
);
readDf.select("c_custkey", "c_name", "c_phone").show(10);
}
}
Maven の pom.xml ファイルに次の依存関係を追加します。
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.13</artifactId>
<version>3.5.4</version>
<scope>provided</scope>
</dependency>
Python
readDf = spark.read.format("hologres").option(
"username", "***").option(
"password", "***").option(
"jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_db").option(
"table", "customer_holo_table").option(
"read.max_task_count", "80").load().filter("c_custkey < 500")
readDf.select("c_custkey", "c_name", "c_phone").show(10)
各プログラミング言語での Spark ジョブの実行:
-
Scala
-
サンプルコードで
sparkselect.scalaファイルを作成し、次のコマンドでジョブを実行します。-- 依存関係をロードします。 spark-shell --jars hologres-connector-spark-3.x-1.5.2-jar-with-dependencies.jar -- ローカルテストの場合、絶対パスでファイルをロードします。 scala> :load D:/sparkselect.scala -
または、依存関係をロードした後、サンプルコードをシェルに直接貼り付けて実行することもできます。
-
-
Java
開発ツールでサンプルコードをインポートし、Maven ツールでパッケージ化します。たとえば、パッケージ化された JAR の名前が
spark_select.jarの場合、次のコマンドでジョブを実行します。-- ジョブJARには絶対パスを使用します。 spark-submit --class SparkSelect --jars hologres-connector-spark-3.x-1.5.2-jar-with-dependencies.jar D:\spark_select.jar -
Python
次のコマンドを実行した後、サンプルコードをシェルに直接貼り付けて実行します。
pyspark --jars hologres-connector-spark-3.x-1.5.2-jar-with-dependencies.jar
パラメータ
一般的なパラメータ
|
パラメータ |
デフォルト |
必須 |
説明 |
|
username |
None |
はい |
|
|
password |
None |
はい |
|
|
table |
None |
はい |
読み取り、または書き込みを行う Hologres テーブルの名前です。 説明
データを読み取る際、代わりに |
|
jdbcurl |
None |
はい |
|
|
enable_serverless_computing |
false |
いいえ |
サーバーレスコンピューティングリソースを使用するかどうかを指定します。このパラメータは、読み取り操作と |
|
serverless_computing_query_priority |
3 |
いいえ |
サーバーレスコンピューティングの実行優先度です。 |
|
statement_timeout_seconds |
28800 (8時間) |
いいえ |
クエリ実行のタイムアウト時間 (秒) です。 |
|
retry_count |
3 |
いいえ |
接続が失敗した場合の再試行回数です。 |
|
direct_connect |
サポートされている場合は、デフォルトで直接接続を使用します。 |
いいえ |
Hologres フロントエンド (アクセスノード) に直接接続するかどうかを指定します。エンドポイントのネットワークスループットは、バッチデータ操作のボトルネックになることがよくあります。デフォルトでは、コネクタはスループットを向上させるために、利用可能な場合は直接接続を使用します。この動作を無効にするには、 |
書き込みパラメーター
Hologres コネクタは、Spark の SaveMode パラメーターをサポートしています。SQL の場合、これは INSERT INTO または INSERT OVERWRITE に対応します。DataFrame の場合、データを書き込む際に SaveMode を追加または上書きに設定できます。上書きモードは、書き込み操作用の一時テーブルを作成し、成功時に元のテーブルを置き換えます。このモードは必要な場合にのみ使用してください。
|
パラメーター |
以前の名前 |
デフォルト |
必須 |
説明 |
|
write.mode |
copy_write_mode |
auto |
いいえ |
書き込みモードです。各書き込みモードの比較については、「Batch write modes」をご参照ください。有効な値は次のとおりです。
|
|
write.copy.max_buffer_size |
max_cell_buffer_size |
52428800 (50 MB) |
いいえ |
|
|
write.copy.dirty_data_check |
copy_write_dirty_data_check |
false |
いいえ |
ダーティデータをチェックするかどうかを指定します。有効にすると、この機能は書き込みに失敗した正確な行を特定できます。ただし、これは書き込みパフォーマンスに影響します。トラブルシューティング以外では、この機能を無効のままにしてください。 |
|
write.on_conflict_action |
なし |
INSERT_OR_REPLACE |
いいえ |
書き込み操作が宛先テーブルでプライマリキーの競合に遭遇した場合に実行するアクション。
|
|
write.stage.compression |
なし |
false |
いいえ |
|
以下のパラメーターは、write.mode が insert に設定されている場合にのみ有効になります。
|
パラメーター |
以前の名前 |
デフォルト |
必須 |
説明 |
|
write.insert.dynamic_partition |
dynamic_partition |
false |
いいえ |
|
|
write.insert.batch_size |
write_batch_size |
512 |
いいえ |
各書き込みスレッドの最大バッチサイズ。 |
|
write.insert.batch_byte_size |
write_batch_byte_size |
2097152 (2 MB) |
いいえ |
各書き込みスレッドの最大バッチサイズ (バイト単位)。デフォルト値は 2 MB です。 |
|
write.insert.max_interval_ms |
write_max_interval_ms |
10000 |
いいえ |
最後のコミットからの経過時間がこの値を超えると、バッチコミットがトリガーされます。 |
|
write.insert.thread_size |
write_thread_size |
1 |
いいえ |
同時書き込みスレッドの数。各スレッドは 1 つのデータベース接続を使用します。 |
|
write.rps_limit |
なし |
-1 |
いいえ |
タスクごとの書き込みのレート制限、秒間行数 (RPS)。デフォルト値の -1 は制限なしを示します。 |
読み取りパラメーター
|
パラメーター |
以前の名前 (v1.5.0 以前) |
デフォルト |
必須 |
説明 |
|
read.mode |
bulk_read |
auto |
いいえ |
読み取りモード。有効な値は次のとおりです:
|
|
read.max_task_count |
max_partition_count |
80 |
いいえ |
データ読み取りのための同時タスクの最大数を指定します。コネクターはテーブルを複数のパーティションに分割し、各パーティションを単一の Spark タスクで処理します。テーブルのシャード数がこの値より小さい場合、パーティションの数はシャード数に制限されます。 |
|
read.copy.max_buffer_size |
/ |
52428800 (50 MB) |
いいえ |
|
|
read.push_down_predicate |
push_down_predicate |
true |
いいえ |
述語プッシュダウンを有効にするかどうかを指定します。有効にすると、フィルター条件や列プルーニングなどの操作をデータソースにプッシュダウンします。 |
|
read.push_down_limit |
push_down_limit |
true |
いいえ |
|
|
read.select.batch_size |
scan_batch_size |
256 |
いいえ |
|
|
read.select.timeout_seconds |
scan_timeout_seconds |
60 |
いいえ |
|
|
read.query |
query |
None |
いいえ |
指定した 説明
|
|
read.split.strategy |
None |
auto |
いいえ |
並列読み取りのためにテーブルデータを複数のタスクに分割する戦略を定義します。有効な値は次のとおりです:
説明
コネクター v1.6.1 以降でサポート。 |
|
read.split.column |
None |
None |
いいえ |
分割列の名前。このパラメーターは、 |
|
read.split.lower_bound |
None |
None |
いいえ |
範囲分割の下限値。このパラメーターは、 |
|
read.split.upper_bound |
None |
None |
いいえ |
範囲分割の上限値。このパラメーターは、 |
|
read.split.num |
None |
None |
いいえ |
読み取り時にデータを分割するタスクの数。このパラメーターは、 |
データ型マッピング
|
Spark 型 |
Hologres 型 |
|
ShortType |
SMALLINT |
|
IntegerType |
INT |
|
LongType |
BIGINT |
|
StringType |
TEXT |
|
StringType |
JSON |
|
StringType |
JSONB |
|
DecimalType |
NUMERIC(38, 18) |
|
BooleanType |
BOOL |
|
DoubleType |
DOUBLE PRECISION |
|
FloatType |
FLOAT |
|
TimestampType |
TIMESTAMPTZ |
|
DateType |
DATE |
|
BinaryType |
BYTEA |
|
BinaryType |
ROARINGBITMAP |
|
ArrayType(IntegerType) |
INT4[] |
|
ArrayType(LongType) |
INT8[] |
|
ArrayType(FloatType) |
FLOAT4[] |
|
ArrayType(DoubleType) |
FLOAT8[] |
|
ArrayType(BooleanType) |
BOOLEAN[] |
|
ArrayType(StringType) |
TEXT[] |
接続数の計算
Hologres-Connector-Spark は、読み取りおよび書き込み操作に一定数の JDBC 接続を使用します。接続数は、次の要因によって決まります:
-
Spark の並列度、つまりジョブ実行中に動作する同時実行タスクの数。これは Spark UI で確認できます。
-
タスクごとに使用される接続数:
-
COPY モードでの書き込み時、各タスクは 1 つの JDBC 接続を使用します。
-
INSERT モードでの書き込み時、各タスクは
write_thread_size個の JDBC 接続を使用します。 -
データ読み取り時、各タスクは 1 つの JDBC 接続を使用します。
-
-
その他の操作:ジョブは、開始時のスキーマ取得などのタスクで、一時的に 1 つの接続を使用することがあります。
ジョブの合計接続数は、次の式で計算できます:
|
項目 |
接続数 |
|
Querying metadata from the catalog |
1 |
|
データの読み取り |
並列度 * 1 + 1 |
|
COPY モードでの書き込み |
並列度 * 1 + 1 |
|
INSERT モードでの書き込み |
並列度 * write_thread_size + 1 |
この計算は、Spark の同時実行タスクのキャパシティが、ジョブが生成するタスクの数を超えていることを前提としています。
Spark の同時実行タスクのキャパシティは、spark.executor.instances などのユーザー設定パラメーターや、Hadoop のファイルブロック分割ポリシーに依存します。詳細については、「Apache Hadoop」をご参照ください。
Spark コネクタのリリースノート
|
バージョン |
リリース日 |
新機能 |
バグ修正 |
|
1.6.1 |
2026-02 |
|
|
|
1.6.0 |
2025-12 |
|
|
|
1.5.6 |
2025-11 |
|
|
|
1.5.5 |
2025-10 |
|
|
|
1.5.4 |
2025-09 |
|
|
|
1.5.2 |
2025-08 |
|
|
|
1.5.0 |
2025-06 |
|
|
|
1.4.2 |
2025-04 |
|
|
|
1.3.2 |
2025-02 |
|
|
|
1.3.0 |
2025-01 |
|
|