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

Hologres:Spark による Hologres の読み取りと書き込み

最終更新日:Sep 12, 2026

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

  1. 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
  2. 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 コネクタのパラメーターの詳細については、「パラメーター」をご参照ください。

  1. 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
  2. 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

はい

  • ご使用の Alibaba Cloud アカウントの AccessKey ID です。[AccessKey Management] ページで取得できます。

  • または、ユーザーアカウント名です。アカウントを作成した後、必要なデータベース権限を付与します。詳細については、「Hologres 権限モデル」をご参照ください。その後、HoloWeb に接続してクエリを実行し、権限を検証します。

password

None

はい

  • ご使用の AccessKey ID に対応する AccessKey Secret です。詳細については、「AccessKey の作成」をご参照ください。

  • または、アカウントのパスワードです。

table

None

はい

読み取り、または書き込みを行う Hologres テーブルの名前です。

説明

データを読み取る際、代わりに read.query パラメータを使用できます。

jdbcurl

None

はい

jdbc:postgresql://:/ 形式の Hologres リアルタイムデータ API の JDBC URL です。これらの詳細を確認するには、Hologres コンソールに移動し、左側メニューで インスタンス一覧 をクリックし、対象のインスタンスを選択して、インスタンスの詳細 ページの ネットワーク情報 セクションでホストとポート番号を確認します。

enable_serverless_computing

false

いいえ

サーバーレスコンピューティングリソースを使用するかどうかを指定します。このパラメータは、読み取り操作と bulk_load モードでの書き込み操作にのみ適用されます。詳細については、「サーバーレスコンピューティングユーザーガイド」をご参照ください。

serverless_computing_query_priority

3

いいえ

サーバーレスコンピューティングの実行優先度です。

statement_timeout_seconds

28800 (8時間)

いいえ

クエリ実行のタイムアウト時間 (秒) です。

retry_count

3

いいえ

接続が失敗した場合の再試行回数です。

direct_connect

サポートされている場合は、デフォルトで直接接続を使用します。

いいえ

Hologres フロントエンド (アクセスノード) に直接接続するかどうかを指定します。エンドポイントのネットワークスループットは、バッチデータ操作のボトルネックになることがよくあります。デフォルトでは、コネクタはスループットを向上させるために、利用可能な場合は直接接続を使用します。この動作を無効にするには、false に設定します。

書き込みパラメーター

Hologres コネクタは、Spark の SaveMode パラメーターをサポートしています。SQL の場合、これは INSERT INTO または INSERT OVERWRITE に対応します。DataFrame の場合、データを書き込む際に SaveMode を追加または上書きに設定できます。上書きモードは、書き込み操作用の一時テーブルを作成し、成功時に元のテーブルを置き換えます。このモードは必要な場合にのみ使用してください。

パラメーター

以前の名前

デフォルト

必須

説明

write.mode

copy_write_mode

auto

いいえ

書き込みモードです。各書き込みモードの比較については、「Batch write modes」をご参照ください。有効な値は次のとおりです。

  • auto (デフォルト):コネクタは、Hologres バージョンと宛先テーブルのメタデータに基づいて最適なモードを自動的に選択します。選択ロジックは次のとおりです。

    1. Hologres インスタンスが v2.2.25 以降で、テーブルにプライマリキーがある場合、bulk_load_on_conflict モードが使用されます。

    2. Hologres インスタンスが v2.1.0 以降で、テーブルにプライマリキーがない場合、bulk_load モードが使用されます。

    3. Hologres インスタンスが v1.3 以降の場合、stream モードが使用されます。

    4. その他の場合、insert モードが使用されます。

  • stream:Fixed Plan を使用して SQL 実行を高速化します。Fixed Plan では、COPY は Hologres v1.3 で導入された機能です。INSERT メソッドと比較して、COPY を使用すると、(ストリーミングの性質により) スループットが向上し、データレイテンシーが低下し、(データをバッチ処理しないため) クライアントのメモリ消費量も削減されます。

    説明

    Hologres コネクタ v1.3.0 以降および Hologres v1.3.34 以降が必要です。

  • bulk_load:バッチ COPY を使用します。Fixed Plan のストリーミング COPY と比較して、バッチ COPY は高 RPS 条件下で Hologres インスタンスへの負荷が低くなりますが、プライマリキーのないテーブルへの書き込みのみをサポートします。

    説明

    Hologres コネクタ v1.4.2 以降および Hologres v2.1.0 以降が必要です。

  • bulk_load_on_conflict:バッチ COPY を使用してプライマリキーを持つテーブルに書き込み、重複するプライマリキーを処理できます。デフォルトでは、プライマリキーを持つ Hologres テーブルへのバッチデータインポートはテーブルロックをトリガーし、複数の接続からの同時書き込み操作が制限されます。コネクタは、宛先テーブルの分散キーに基づいてデータを再分散させ、各 Spark タスクが単一のシャードに書き込めるようにします。これにより、テーブルロックがシャードレベルのロックに削減され、同時書き込みが可能になり、書き込みパフォーマンスが向上します。各接続は少数のシャードのデータのみを保持するため、この最適化は、小さなファイルの数を大幅に削減し、Hologres のメモリ使用量を低下させます。テストでは、同時書き込み前にデータを再パーティション化すると、stream モードと比較してシステム負荷が約 67% 削減されることがわかっています。

    説明

    Hologres コネクタ v1.4.2 以降および Hologres v2.2.25 以降が必要です。

  • insert:INSERT メソッドを使用してデータを書き込みます。

  • stage:ステージテーブルを使用してバッチ書き込みを行います。このモードは、一時ストレージステージを通じてニアリアルタイムでデータをインポートし、INSERT FROM SELECT ステートメントを使用して宛先テーブルにロードします。バッチシナリオでは、このモードは他のモードよりもパフォーマンスが優れており、システム負荷も低くなります。また、サーバーレスコンピューティングリソースの使用もサポートしています。INSERT INTO と INSERT OVERWRITE の両方の操作をサポートします。Hologres v4.1.0 以降およびコネクタ v1.6.1 以降が必要です。

write.copy.max_buffer_size

max_cell_buffer_size

52428800 (50 MB)

いいえ

COPY モードで書き込む際のローカルバッファーの最大サイズ。通常、この値を調整する必要はありません。ただし、非常に長い文字列などの大きなフィールドを書き込む際にバッファーオーバーフローが発生する場合は、この値を増やすことができます。

write.copy.dirty_data_check

copy_write_dirty_data_check

false

いいえ

ダーティデータをチェックするかどうかを指定します。有効にすると、この機能は書き込みに失敗した正確な行を特定できます。ただし、これは書き込みパフォーマンスに影響します。トラブルシューティング以外では、この機能を無効のままにしてください。

write.on_conflict_action

なし

INSERT_OR_REPLACE

いいえ

書き込み操作が宛先テーブルでプライマリキーの競合に遭遇した場合に実行するアクション。

  • INSERT_OR_IGNORE:プライマリキーの競合が発生した場合、データは書き込まれません。

  • INSERT_OR_UPDATE:プライマリキーの競合が発生した場合、指定された列が更新されます。

  • INSERT_OR_REPLACE:プライマリキーの競合が発生した場合、すべての列が更新されます。

write.stage.compression

なし

false

いいえ

write.mode が stage に設定されている場合に有効になります。これを true に設定すると、ステージ書き込み中の圧縮が有効になります。Hologres v4.2.8 以降およびコネクタ v1.6.2 以降が必要です。

以下のパラメーターは、write.mode が insert に設定されている場合にのみ有効になります。

パラメーター

以前の名前

デフォルト

必須

説明

write.insert.dynamic_partition

dynamic_partition

false

いいえ

write.mode が insert の場合、これを true に設定すると、親パーティションテーブルに書き込む際に存在しないパーティションが自動的に作成されます。

write.insert.batch_size

write_batch_size

512

いいえ

各書き込みスレッドの最大バッチサイズ。Put 操作の数がこの値に達すると、バッチコミットがトリガーされます。

write.insert.batch_byte_size

write_batch_byte_size

2097152 (2 MB)

いいえ

各書き込みスレッドの最大バッチサイズ (バイト単位)。デフォルト値は 2 MB です。Put データのバイトサイズがこの値に達すると、バッチコミットがトリガーされます。

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

いいえ

読み取りモード。有効な値は次のとおりです:

  • auto (デフォルト):コネクターは、Hologres のバージョンとテーブルメタデータに基づいて最適なモードを自動的に選択します。選択ロジックは次のとおりです:

    1. 読み取るフィールドに JSONB データ型が含まれる場合、select モードを使用します。

    2. インスタンスが v3.0.24 以降の場合、bulk_read_compressed モードを使用します。

    3. その他の場合は、bulk_read モードを使用します。

  • bulk_read:COPY OUT メソッドを使用して Arrow フォーマットでデータを読み取ります。これは select モードよりも数倍高速です。現在、Hologres から JSONB データ型を読み取ることはサポートしていません。

  • bulk_read_compressed:COPY OUT メソッドを使用して、Arrow フォーマットで圧縮データを読み取ります。これにより、非圧縮データを読み取る場合と比較して帯域幅を約 45% 節約できます。

  • select:標準の SELECT ステートメントを使用してデータを読み取ります。

read.max_task_count

max_partition_count

80

いいえ

データ読み取りのための同時タスクの最大数を指定します。コネクターはテーブルを複数のパーティションに分割し、各パーティションを単一の Spark タスクで処理します。テーブルのシャード数がこの値より小さい場合、パーティションの数はシャード数に制限されます。

read.copy.max_buffer_size

/

52428800 (50 MB)

いいえ

COPY モードで読み取る際のローカルバッファーの最大サイズ。フィールドサイズが大きいために例外が発生した場合は、この値を増やしてください。

read.push_down_predicate

push_down_predicate

true

いいえ

述語プッシュダウンを有効にするかどうかを指定します。有効にすると、フィルター条件や列プルーニングなどの操作をデータソースにプッシュダウンします。

read.push_down_limit

push_down_limit

true

いいえ

LIMIT プッシュダウンを有効にするかどうかを指定します。

read.select.batch_size

scan_batch_size

256

いいえ

read.mode が select に設定されている場合に有効。Hologres から読み取る際に、1 回のスキャン操作でフェッチする行数を指定します。

read.select.timeout_seconds

scan_timeout_seconds

60

いいえ

read.mode が select に設定されている場合に有効。Hologres から読み取る際のスキャン操作のタイムアウト期間を指定します。

read.query

query

None

いいえ

指定した query を使用して Hologres から読み取ります。このパラメーターは table パラメーターと併用できません。

説明
  • query メソッドを使用してデータを読み取る場合、タスクは1つしか使用できません。述語プッシュダウンはサポートしていません。

  • table メソッドを使用してデータを読み取る場合、読み取り操作は Hologres テーブルのシャード数に基づいて並列実行される複数のタスクに分割されます。

read.split.strategy

None

auto

いいえ

並列読み取りのためにテーブルデータを複数のタスクに分割する戦略を定義します。有効な値は次のとおりです:

  • auto (デフォルト):コネクターが最適な戦略を自動的に選択します。

  • shard:Hologres のシャードに基づいてデータをシャーディングし、各タスクがシャードの一部を読み取ります。これは、シャードの分散が明確なテーブルに適しています。

  • range:指定した列の値の範囲に基づいてデータをシャーディングします。この設定には、read.split.column、read.split.lower_bound、read.split.upper_bound、および read.split.num パラメーターが必要です。

  • partition:パーティション列の値に基づいてデータをシャーディングします。このオプションは、read.split.column および read.split.num パラメーターとともに使用する必要があります。パーティションテーブルまたはパーティションビューに適用されます。

説明

コネクター v1.6.1 以降でサポート。

read.split.column

None

None

いいえ

分割列の名前。このパラメーターは、read.split.strategy が range または partition の場合に必要。

read.split.lower_bound

None

None

いいえ

範囲分割の下限値。このパラメーターは、read.split.strategy が range に設定されている場合に必要。

read.split.upper_bound

None

None

いいえ

範囲分割の上限値。このパラメーターは、read.split.strategy が range に設定されている場合に必要。

read.split.num

None

None

いいえ

読み取り時にデータを分割するタスクの数。このパラメーターは、read.split.strategy が range または partition の場合に必要。

データ型マッピング

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

  • ステージ書き込みモードを追加しました。このモードは一時的なステージテーブルを使用して書き込みパフォーマンスを向上させます。この機能には Hologres v4.1.0 以降が必要です。

  • ステージモードで INSERT INTO および INSERT OVERWRITE 操作をサポートしました。

  • overwrite 操作における一時テーブルのクリーンアップロジックを最適化しました。

  • タスクごとの書き込みレートを制御する write.rps_limit パラメーターを追加しました。

  • overwrite 操作中に DDL を実行した際に発生する強制同期リプレイの問題を修正しました。

1.6.0

2025-12

  • read.split.strategy パラメーターを追加し、シャード、レンジ、パーティション の 3 つの読み取り分割戦略をサポートしました。

  • Arrow ベースの読み取りで JSONB データ型をサポートしました。

  • write.insert.ignore_null_when_update パラメーターを追加し、UPDATE 操作中に NULL 値を無視できるようになりました。

  • JaCoCo テストカバレッジレポートをサポートしました。

  • 1970 年より前の date 値を書き込む際のエラーを修正しました。

1.5.6

2025-11

  • AKV4 認証方式をサポートしました。

  • write.copy.disable_right_join パラメーターを追加し、COPY の書き込みパフォーマンスを最適化しました。

  • カタログ書き込みにおける列型のチェックロジックを最適化しました。

1.5.5

2025-10

  • 低精度のデータを高精度のデータ型へ書き込めるようになりました。

  • トラブルシューティングを簡素化するために、ログに appname や taskid などのプレフィックスを追加しました。

  • COPY OUT 操作で Arrow LZ4 圧縮を有効にしました。

  • パラメーターのフォーマットと検証を最適化しました。

1.5.4

2025-09

  • remove_u0000 パラメーターを追加し、テキスト列から U+0000 文字を自動的に削除できるようになりました。

  • Hologres カタログをリファクタリングし、名前空間をスキーマにマッピングするようにしました。

  • データシャーディング用のユーティリティクラス RepartitionUtil を追加しました。

  • ドキュメントをリファクタリングし、リンクを更新しました。

1.5.2

2025-08

  • Hologres カタログを導入し、外部テーブルによる Hologres へのデータの読み書きが可能になりました。

  • 述語プッシュダウンとリミットプッシュダウンに対応しました。

  • 読み取り操作で bulk_read バッチモードを有効にし、パフォーマンスを桁違いに向上させました。

  • 読み取りパフォーマンスを最適化しました。

1.5.0

2025-06

  • 述語プッシュダウン、LIMIT プッシュダウン、カラムプルーニングをサポートしました。

  • SELECT クエリを使用してデータを読み取れるようになりました。

  • 読み取り操作でバッチモードを有効にしました。

  • 読み取り関連の問題を修正しました。

1.4.2

2025-04

  • bulk_load 書き込みモードを導入し、プライマリキーのないテーブルへのバッチ書き込みが可能になりました。

  • bulk_load_on_conflict モードを追加し、プライマリキーのあるテーブルでのプライマリキー競合を処理できるようになりました。

  • 書き込みパフォーマンスを最適化しました。

1.3.2

2025-02

  • Hologres からのデータ読み取りをサポートしました。

  • シャードによる並列読み取りを有効にしました。

  • 自動スキーマ推論をサポートしました。

  • 初期の読み取り機能を安定させました。

1.3.0

2025-01

  • stream (固定 COPY) 書き込みモードを導入しました。

  • Hologres v1.3 の固定 COPY 機能をサポートしました。

  • 書き込みスループットとレイテンシーを最適化しました。