Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Spark JDBC

ClickHouse対応

JDBC は、Spark で最もよく使われているデータソースの 1 つです。 このセクションでは、Spark で ClickHouse 公式 JDBC コネクタ を 使用する方法について詳しく説明します。

データの読み取り

public static void main(String[] args) {
        // Sparkセッションを初期化
        SparkSession spark = SparkSession.builder().appName("example").master("local").getOrCreate();

        String jdbcURL = "jdbc:ch://localhost:8123/default";
        String query = "select * from example_table where id > 2";

        //---------------------------------------------------------------------------------------------------
        // jdbcメソッドを使用してClickHouseからテーブルを読み込む
        //---------------------------------------------------------------------------------------------------
        Properties jdbcProperties = new Properties();
        jdbcProperties.put("user", "default");
        jdbcProperties.put("password", "123456");

        Dataset<Row> df1 = spark.read().jdbc(jdbcURL, String.format("(%s)", query), jdbcProperties);

        df1.show();

        //---------------------------------------------------------------------------------------------------
        // loadメソッドを使用してClickHouseからテーブルを読み込む
        //---------------------------------------------------------------------------------------------------
        Dataset<Row> df2 = spark.read()
                .format("jdbc")
                .option("url", jdbcURL)
                .option("user", "default")
                .option("password", "123456")
                .option("query", query)
                .load();

        df2.show();

        // Sparkセッションを停止
        spark.stop();
    }

データの書き込み

 public static void main(String[] args) {
        // Sparkセッションを初期化する
        SparkSession spark = SparkSession.builder().appName("example").master("local").getOrCreate();

        // JDBC接続情報
        String jdbcUrl = "jdbc:ch://localhost:8123/default";
        Properties jdbcProperties = new Properties();
        jdbcProperties.put("user", "default");
        jdbcProperties.put("password", "123456");

        // サンプルDataFrameを作成する
        StructType schema = new StructType(new StructField[]{
                DataTypes.createStructField("id", DataTypes.IntegerType, false),
                DataTypes.createStructField("name", DataTypes.StringType, false)
        });

        List<Row> rows = new ArrayList<Row>();
        rows.add(RowFactory.create(1, "John"));
        rows.add(RowFactory.create(2, "Doe"));

        Dataset<Row> df = spark.createDataFrame(rows, schema);

        //---------------------------------------------------------------------------------------------------
        // jdbcメソッドを使用してdfをClickHouseに書き込む
        //---------------------------------------------------------------------------------------------------

        df.write()
                .mode(SaveMode.Append)
                .jdbc(jdbcUrl, "example_table", jdbcProperties);

        //---------------------------------------------------------------------------------------------------
        // saveメソッドを使用してdfをClickHouseに書き込む
        //---------------------------------------------------------------------------------------------------

        df.write()
                .format("jdbc")
                .mode("append")
                .option("url", jdbcUrl)
                .option("dbtable", "example_table")
                .option("user", "default")
                .option("password", "123456")
                .save();

        // Sparkセッションを停止する
        spark.stop();
    }

並列度

Spark JDBC を使用する場合、Spark は単一のパーティションでデータを読み取ります。より高い同時実行性を実現するには、 partitionColumnlowerBoundupperBoundnumPartitions を指定する必要があります。これらの設定は、複数のワーカーから 並列に読み取る際に、テーブルをどのようにパーティション分割するかを定義します。 詳しくは、Apache Spark の公式ドキュメントの JDBC 構成 を参照してください。

JDBC の制限事項

  • ClickHouse dialect がないため、Spark JDBC は複合型 (MAP、ARRAY、STRUCT) をサポートしていません。複合型を完全にサポートするには、ネイティブの Spark-ClickHouse コネクタを使用してください。
  • 現時点では、JDBC を使用してデータを挿入できるのは既存のテーブルのみです (現在、Spark が他のコネクタで行うように、DF 挿入時にテーブルを自動作成する方法はありません) 。
Navigation