Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Spark 커넥터

ClickHouse 지원

이 커넥터는 고급 파티셔닝과 프레디케이트 푸시다운 같은 ClickHouse 전용 최적화를 활용해 쿼리 성능과 데이터 처리 효율을 향상시킵니다. 이 커넥터는 ClickHouse의 공식 JDBC 커넥터를 기반으로 하며, 자체 카탈로그를 관리합니다.

Spark 3.0 이전에는 Spark에 내장 카탈로그 개념이 없었기 때문에, 일반적으로 Hive Metastore나 AWS Glue 같은 외부 카탈로그 시스템에 의존했습니다. 이러한 외부 솔루션에서는 Spark에서 사용하기 전에 데이터 소스 테이블을 수동으로 등록해야 했습니다. 하지만 Spark 3.0에 카탈로그 개념이 도입되면서 Spark는 카탈로그 플러그인을 등록해 테이블을 자동으로 검색할 수 있게 되었습니다.

Spark의 기본 카탈로그는 spark_catalog이며, 테이블은 {catalog name}.{database}.{table} 형식으로 식별됩니다. 새로운 카탈로그 기능을 사용하면 이제 단일 Spark 애플리케이션에서 여러 카탈로그를 추가해 사용할 수 있습니다.

Catalog API와 TableProvider API 중 선택하기

ClickHouse Spark connector는 Catalog APITableProvider API(포맷 기반 접근 방식)라는 두 가지 액세스 패턴을 지원합니다. 두 방식의 차이를 이해하면 사용 사례에 맞는 접근 방식을 선택하는 데 도움이 됩니다.

Catalog API와 TableProvider API 비교

기능 Catalog API TableProvider API
구성 Spark 구성으로 중앙 집중식 관리 작업별 옵션으로 지정
테이블 검색 카탈로그를 통해 자동으로 검색 테이블을 수동으로 지정
DDL 작업 완전 지원 (CREATE, DROP, ALTER) 제한적 지원 (자동 테이블 생성만 가능)
Spark SQL 통합 네이티브 (clickhouse.database.table) 포맷 지정 필요
사용 사례 중앙 집중식 구성을 사용하는 장기적이고 안정적인 연결 애드혹, 동적 또는 임시 접근

요구 사항

  • Java 8 또는 17 (Spark 4.0에는 Java 17 이상이 필요합니다)
  • Scala 2.12 또는 2.13 (Spark 4.0은 Scala 2.13만 지원합니다)
  • Apache Spark 3.3, 3.4, 3.5 또는 4.0

호환성 매트릭스

버전 호환되는 Spark 버전 ClickHouse JDBC 버전
main Spark 3.3, 3.4, 3.5, 4.0 0.9.4
0.10.0 Spark 3.3, 3.4, 3.5, 4.0 0.9.5
0.9.0 Spark 3.3, 3.4, 3.5, 4.0 0.9.4
0.8.1 Spark 3.3, 3.4, 3.5 0.6.3
0.7.3 Spark 3.3, 3.4 0.4.6
0.6.0 Spark 3.3 0.3.2-patch11
0.5.0 Spark 3.2, 3.3 0.3.2-patch11
0.4.0 Spark 3.2, 3.3 의존하지 않음
0.3.0 Spark 3.2, 3.3 의존하지 않음
0.2.1 Spark 3.2 의존하지 않음
0.1.2 Spark 3.2 의존하지 않음

설치 및 설정

ClickHouse를 Spark와 통합하는 방법은 프로젝트 구성에 따라 여러 가지가 있습니다. 프로젝트의 빌드 파일(Maven의 pom.xml 또는 SBT의 build.sbt 등)에 ClickHouse Spark 커넥터를 의존성으로 직접 추가할 수 있습니다. 또는 필요한 JAR 파일을 $SPARK_HOME/jars/ 폴더에 넣거나, spark-submit 명령에서 --jars 플래그를 사용해 Spark 옵션으로 직접 전달할 수 있습니다. 두 방법 모두 Spark 환경에서 ClickHouse 커넥터를 사용할 수 있게 해줍니다.

의존성으로 추가하기

<dependency>
  <groupId>com.clickhouse.spark</groupId>
  <artifactId>clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}</artifactId>
  <version>{{ stable_version }}</version>
</dependency>
<dependency>
  <groupId>com.clickhouse</groupId>
  <artifactId>clickhouse-jdbc</artifactId>
  <classifier>all</classifier>
  <version>{{ clickhouse_jdbc_version }}</version>
  <exclusions>
    <exclusion>
      <groupId>*</groupId>
      <artifactId>*</artifactId>
    </exclusion>
  </exclusions>
</dependency>

SNAPSHOT 버전을 사용하려면 Maven에서 Sonatype의 SNAPSHOT 릴리스 사용 지침을 따르십시오.

라이브러리 다운로드

바이너리 JAR 파일의 이름 패턴은 다음과 같습니다:

clickhouse-spark-runtime-${spark_binary_version}_${scala_binary_version}-${version}.jar

사용 가능한 모든 릴리스 JAR 파일은 Maven Central Repository에서 확인할 수 있습니다. 일별 빌드 SNAPSHOT JAR 파일은 위에서 구성한 Sonatype snapshots 리포지토리를 통해 사용할 수 있습니다.

카탈로그 등록(필수)

ClickHouse 테이블에 액세스하려면 다음 구성으로 새 Spark 카탈로그를 설정해야 합니다.

Property Value Default Value Required
spark.sql.catalog.<catalog_name> com.clickhouse.spark.ClickHouseCatalog N/A Yes
spark.sql.catalog.<catalog_name>.host <clickhouse_host> localhost No
spark.sql.catalog.<catalog_name>.protocol http http No
spark.sql.catalog.<catalog_name>.http_port <clickhouse_port> 8123 No
spark.sql.catalog.<catalog_name>.user <clickhouse_username> default No
spark.sql.catalog.<catalog_name>.password <clickhouse_password> (빈 문자열) No
spark.sql.catalog.<catalog_name>.database <database> default No
spark.<catalog_name>.write.format json arrow No

이 설정은 다음 방법 중 하나로 지정할 수 있습니다.

  • spark-defaults.conf를 편집하거나 생성합니다.
  • 구성을 spark-submit 명령에 전달합니다(또는 spark-shell/spark-sql CLI 명령에 전달).
  • Context를 초기화할 때 구성을 추가합니다.

TableProvider API 사용하기 (포맷 기반 접근 방식)

카탈로그 기반 접근 방식 외에도, ClickHouse Spark 커넥터는 TableProvider API를 통한 포맷 기반 접근 방식을 지원합니다.

포맷 기반 읽기 예시

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

# format API를 사용해 ClickHouse에서 읽기
df = spark.read \
    .format("clickhouse") \
    .option("host", "your-clickhouse-host") \
    .option("protocol", "https") \
    .option("http_port", "8443") \
    .option("database", "default") \
    .option("table", "your_table") \
    .option("user", "default") \
    .option("password", "your_password") \
    .option("ssl", "true") \
    .load()

df.show()

포맷 기반 쓰기 예시

# 포맷 API를 사용해 ClickHouse에 쓰기
df.write \
    .format("clickhouse") \
    .option("host", "your-clickhouse-host") \
    .option("protocol", "https") \
    .option("http_port", "8443") \
    .option("database", "default") \
    .option("table", "your_table") \
    .option("user", "default") \
    .option("password", "your_password") \
    .option("ssl", "true") \
    .mode("append") \
    .save()

TableProvider 기능

TableProvider API는 여러 강력한 기능을 제공합니다:

자동 테이블 생성

존재하지 않는 테이블에 쓸 경우 커넥터가 적절한 스키마로 테이블을 자동 생성합니다. 커넥터는 다음과 같은 합리적인 기본값을 제공합니다.

  • Engine: 지정하지 않으면 기본값으로 MergeTree()를 사용합니다. engine 옵션을 사용해 다른 엔진을 지정할 수 있습니다(예: ReplacingMergeTree(), SummingMergeTree() 등).
  • ORDER BY: 필수 - 새 테이블을 생성할 때는 반드시 order_by 옵션을 명시적으로 지정해야 합니다. 커넥터는 지정된 모든 컬럼이 스키마에 존재하는지 검증합니다.
  • 널 허용 키 지원: ORDER BY에 널 허용 컬럼이 포함되어 있으면 settings.allow_nullable_key=1을 자동으로 추가합니다
# 명시적으로 ORDER BY를 지정하면 테이블이 자동 생성됩니다(필수)
df.write \
    .format("clickhouse") \
    .option("host", "your-host") \
    .option("database", "default") \
    .option("table", "new_table") \
    .option("order_by", "id") \
    .mode("append") \
    .save()

# 사용자 지정 엔진으로 테이블 생성 옵션 지정
df.write \
    .format("clickhouse") \
    .option("host", "your-host") \
    .option("database", "default") \
    .option("table", "new_table") \
    .option("order_by", "id, timestamp") \
    .option("engine", "ReplacingMergeTree()") \
    .option("settings.allow_nullable_key", "1") \
    .mode("append") \
    .save()

TableProvider 연결 옵션

포맷 기반 API를 사용할 때는 다음 연결 옵션을 사용할 수 있습니다:

연결 옵션

Option Description Default Value Required
host ClickHouse 서버 호스트명 localhost
protocol 연결 프로토콜(http 또는 https) http 아니요
http_port HTTP/HTTPS 포트 8123 아니요
database 데이터베이스 이름 default
table 테이블 이름 N/A
user authentication용 사용자 이름 default 아니요
password authentication용 비밀번호 (빈 문자열) 아니요
ssl SSL 연결 사용 여부 false 아니요
ssl_mode SSL 모드(NONE, STRICT 등) STRICT 아니요
timezone 날짜/시간 작업에 사용할 시간대 server 아니요

테이블 생성 옵션

이 옵션은 테이블이 아직 없어서 새로 생성해야 할 때 사용됩니다:

Option Description Default Value Required
order_by ORDER BY 절에 사용할 컬럼입니다. 여러 컬럼은 쉼표로 구분합니다. N/A
engine ClickHouse 테이블 엔진(예: MergeTree(), ReplacingMergeTree(), SummingMergeTree() 등) MergeTree() 아니요
settings.allow_nullable_key ORDER BY에서 널 허용 키를 활성화합니다(ClickHouse Cloud용). 자동 감지** 아니요
settings.<key> 임의의 ClickHouse 테이블 설정 N/A 아니요
cluster 분산 테이블용 클러스터 이름 N/A 아니요
clickhouse.column.<name>.variant_types Variant 컬럼에 사용할 ClickHouse 타입의 쉼표로 구분된 목록입니다(예: String, Int64, Bool, JSON). 타입 이름은 대소문자를 구분합니다. 쉼표 뒤 공백은 있어도 되고 없어도 됩니다. N/A 아니요
  • 새 테이블을 생성할 때는 order_by 옵션이 필요합니다. 지정한 모든 컬럼은 스키마에 존재해야 합니다. ** ORDER BY에 널 허용 컬럼이 포함되어 있고 이 값이 명시적으로 지정되지 않으면 자동으로 1로 설정됩니다.

쓰기 모드

Spark 커넥터(TableProvider API와 Catalog API 모두)는 다음 Spark 쓰기 모드를 지원합니다.

  • append: 기존 테이블에 데이터를 추가합니다
  • overwrite: 테이블의 모든 데이터를 대체합니다(테이블을 TRUNCATE함)
# Overwrite 모드(먼저 테이블을 TRUNCATE함)
df.write \
    .format("clickhouse") \
    .option("host", "your-host") \
    .option("database", "default") \
    .option("table", "my_table") \
    .mode("overwrite") \
    .save()

ClickHouse 옵션 구성

Catalog API와 TableProvider API는 모두 ClickHouse 전용 옵션(커넥터 옵션 제외)의 구성을 지원합니다. 이러한 옵션은 테이블을 생성하거나 쿼리를 실행할 때 ClickHouse에 전달됩니다.

ClickHouse 옵션을 사용하면 allow_nullable_key, index_granularity와 같은 ClickHouse 전용 설정과 그 밖의 테이블 수준 또는 쿼리 수준 설정을 구성할 수 있습니다. 이는 커넥터가 ClickHouse에 연결하는 방식을 제어하는 커넥터 옵션(host, database, table 등)과는 다릅니다.

TableProvider API 사용

TableProvider API에서는 settings.<key> 옵션 포맷을 사용합니다:

df.write \
    .format("clickhouse") \
    .option("host", "your-host") \
    .option("database", "default") \
    .option("table", "my_table") \
    .option("order_by", "id") \
    .option("settings.allow_nullable_key", "1") \
    .option("settings.index_granularity", "8192") \
    .mode("append") \
    .save()

Catalog API 사용

Catalog API를 사용할 때는 Spark 구성에서 spark.sql.catalog.<catalog_name>.option.<key> 포맷을 사용하십시오:

spark.sql.catalog.clickhouse.option.allow_nullable_key 1
spark.sql.catalog.clickhouse.option.index_granularity 8192

또는 Spark SQL로 테이블을 생성할 때 설정할 수도 있습니다:

CREATE TABLE clickhouse.default.my_table (
  id INT,
  name STRING
) USING ClickHouse
TBLPROPERTIES (
  engine = 'MergeTree()',
  order_by = 'id',
  'settings.allow_nullable_key' = '1',
  'settings.index_granularity' = '8192'
)

ClickHouse Cloud 설정

ClickHouse Cloud에 연결할 때는 SSL을 활성화하고 적절한 SSL 모드를 설정하십시오. 예시는 다음과 같습니다.

spark.sql.catalog.clickhouse.option.ssl        true
spark.sql.catalog.clickhouse.option.ssl_mode   NONE

데이터 읽기

public static void main(String[] args) {
        // Spark 세션 생성
        SparkSession spark = SparkSession.builder()
                .appName("example")
                .master("local[*]")
                .config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
                .config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
                .config("spark.sql.catalog.clickhouse.protocol", "http")
                .config("spark.sql.catalog.clickhouse.http_port", "8123")
                .config("spark.sql.catalog.clickhouse.user", "default")
                .config("spark.sql.catalog.clickhouse.password", "123456")
                .config("spark.sql.catalog.clickhouse.database", "default")
                .config("spark.clickhouse.write.format", "json")
                .getOrCreate();

        Dataset<Row> df = spark.sql("select * from clickhouse.default.example_table");

        df.show();

        spark.stop();
    }

데이터 쓰기

 public static void main(String[] args) throws AnalysisException {

        // Spark 세션을 생성합니다.
        SparkSession spark = SparkSession.builder()
                .appName("example")
                .master("local[*]")
                .config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
                .config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
                .config("spark.sql.catalog.clickhouse.protocol", "http")
                .config("spark.sql.catalog.clickhouse.http_port", "8123")
                .config("spark.sql.catalog.clickhouse.user", "default")
                .config("spark.sql.catalog.clickhouse.password", "123456")
                .config("spark.sql.catalog.clickhouse.database", "default")
                .config("spark.clickhouse.write.format", "json")
                .getOrCreate();

        // DataFrame의 스키마를 정의합니다.
        StructType schema = new StructType(new StructField[]{
                DataTypes.createStructField("id", DataTypes.IntegerType, false),
                DataTypes.createStructField("name", DataTypes.StringType, false),
        });

        List<Row> data = Arrays.asList(
                RowFactory.create(1, "Alice"),
                RowFactory.create(2, "Bob")
        );

        // DataFrame을 생성합니다.
        Dataset<Row> df = spark.createDataFrame(data, schema);

        df.writeTo("clickhouse.default.example_table").append();

        spark.stop();
    }

DDL 작업

Spark SQL을 사용해 ClickHouse 인스턴스에서 DDL 작업을 수행할 수 있으며, 모든 변경 사항은 즉시 ClickHouse에 저장됩니다. Spark SQL에서는 ClickHouse에서와 동일하게 쿼리를 작성할 수 있으므로, 예를 들어 CREATE TABLE, TRUNCATE 등의 명령을 수정 없이 직접 실행할 수 있습니다:

USE clickhouse; 

CREATE TABLE test_db.tbl_sql (
  create_time TIMESTAMP NOT NULL,
  m           INT       NOT NULL COMMENT 'part key',
  id          BIGINT    NOT NULL COMMENT 'sort key',
  value       STRING
) USING ClickHouse
PARTITIONED BY (m)
TBLPROPERTIES (
  engine = 'MergeTree()',
  order_by = 'id',
  settings.index_granularity = 8192
);

위의 예시는 Spark SQL 쿼리를 보여 주며, Java, Scala, PySpark 또는 셸 등 어떤 API로든 애플리케이션 내에서 실행할 수 있습니다.

VariantType 사용하기

커넥터는 반정형 데이터를 다루기 위해 Spark의 VariantType을 지원합니다. VariantType은 ClickHouse의 JSONVariant 타입에 매핑되므로, 유연한 스키마의 데이터를 효율적으로 저장하고 쿼리할 수 있습니다.

ClickHouse 타입 매핑

ClickHouse 유형 Spark 유형 설명
JSON VariantType JSON 객체만 저장합니다({로 시작해야 합니다)
Variant(T1, T2, ...) VariantType 기본 타입, 배열, JSON을 포함한 여러 타입을 저장합니다

VariantType 데이터 읽기

ClickHouse에서 데이터를 읽으면 JSONVariant 컬럼이 자동으로 Spark의 VariantType에 매핑됩니다.

// JSON 컬럼을 VariantType으로 읽기
val df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")

// variant 데이터에 액세스
df.show()

// 확인할 수 있도록 variant를 JSON 문자열로 변환
import org.apache.spark.sql.functions._
df.select(
  col("id"),
  to_json(col("data")).as("data_json")
).show()

VariantType 데이터 쓰기

JSON 또는 Variant 컬럼 타입을 사용해 VariantType 데이터를 ClickHouse에 쓸 수 있습니다:

import org.apache.spark.sql.functions._

// JSON 데이터로 DataFrame 생성
val jsonData = Seq(
  (1, """{"name": "Alice", "age": 30}"""),
  (2, """{"name": "Bob", "age": 25}"""),
  (3, """{"name": "Charlie", "city": "NYC"}""")
).toDF("id", "json_string")

// JSON 문자열을 VariantType으로 파싱
val variantDF = jsonData.select(
  col("id"),
  parse_json(col("json_string")).as("data")
)

// JSON 타입으로 ClickHouse에 쓰기 (JSON 객체만 지원)
variantDF.writeTo("clickhouse.default.user_data").create()

// 또는 여러 타입을 포함하는 Variant 지정
spark.sql("""
  CREATE TABLE clickhouse.default.mixed_data (
    id INT,
    data VARIANT
  ) USING clickhouse
  TBLPROPERTIES (
    'clickhouse.column.data.variant_types' = 'String, Int64, Bool, JSON',
    'engine' = 'MergeTree()',
    'order_by' = 'id'
  )
""")

Spark SQL로 VariantType 테이블 생성하기

Spark SQL DDL을 사용해 VariantType 테이블을 생성할 수 있습니다:

-- JSON 타입으로 테이블 생성 (기본값)
CREATE TABLE clickhouse.default.json_table (
  id INT,
  data VARIANT
) USING clickhouse
TBLPROPERTIES (
  'engine' = 'MergeTree()',
  'order_by' = 'id'
)
-- 여러 타입을 지원하는 Variant 타입 테이블 생성
CREATE TABLE clickhouse.default.flexible_data (
  id INT,
  data VARIANT
) USING clickhouse
TBLPROPERTIES (
  'clickhouse.column.data.variant_types' = 'String, Int64, Float64, Bool, Array(String), JSON',
  'engine' = 'MergeTree()',
  'order_by' = 'id'
)

Variant 타입 구성하기

VariantType 컬럼이 포함된 테이블을 생성할 때 사용할 ClickHouse 타입을 지정할 수 있습니다:

JSON 타입 (기본값)

variant_types 속성을 지정하지 않으면 해당 컬럼은 기본적으로 ClickHouse의 JSON 타입을 사용하며, 이 타입은 JSON 객체만 허용합니다:

CREATE TABLE clickhouse.default.json_table (
  id INT,
  data VARIANT
) USING clickhouse
TBLPROPERTIES (
  'engine' = 'MergeTree()',
  'order_by' = 'id'
)

다음과 같은 ClickHouse 쿼리가 생성됩니다:

CREATE TABLE json_table (id Int32, data JSON) ENGINE = MergeTree() ORDER BY id

여러 타입을 지원하는 Variant Type

기본 타입, 배열, JSON 객체를 지원하려면 variant_types 속성에 타입을 지정합니다:

CREATE TABLE clickhouse.default.flexible_data (
  id INT,
  data VARIANT
) USING clickhouse
TBLPROPERTIES (
  'clickhouse.column.data.variant_types' = 'String, Int64, Float64, Bool, Array(String), JSON',
  'engine' = 'MergeTree()',
  'order_by' = 'id'
)

다음과 같은 ClickHouse 쿼리가 생성됩니다:

CREATE TABLE flexible_data (
  id Int32, 
  data Variant(String, Int64, Float64, Bool, Array(String), JSON)
) ENGINE = MergeTree() ORDER BY id

지원되는 Variant 타입

다음 ClickHouse 타입은 Variant()에서 사용할 수 있습니다.

  • 기본 타입: String, Int8, Int16, Int32, Int64, UInt8, UInt16, UInt32, UInt64, Float32, Float64, Bool
  • 배열: Array(T) — 여기서 T는 중첩 배열을 포함한 지원되는 모든 타입입니다
  • JSON: JSON 객체 저장용 JSON

읽기 포맷 구성

기본적으로 JSON 및 Variant 컬럼은 VariantType으로 읽힙니다. 이 동작은 재정의할 수 있으며, 문자열로 읽도록 설정할 수 있습니다.

// JSON/Variant를 VariantType 대신 문자열로 읽기
spark.conf.set("spark.clickhouse.read.jsonAs", "string")

val df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
// data 컬럼은 JSON 문자열을 포함하는 StringType으로 읽힘

쓰기 포맷 지원

VariantType의 쓰기 지원은 포맷에 따라 다릅니다:

포맷 지원 참고 사항
JSON ✅ 전체 지원 JSONVariant 타입을 모두 지원합니다. VariantType 데이터에는 JSON 사용을 권장합니다
Arrow ⚠️ 부분 지원 ClickHouse JSON 타입으로 쓰기를 지원합니다. ClickHouse Variant 타입은 지원하지 않습니다. 전체 지원은 https://github.com/ClickHouse/ClickHouse/issues/92752 이슈가 해결되면 제공될 예정입니다

쓰기 포맷을 설정합니다:

spark.conf.set("spark.clickhouse.write.format", "json")  // Variant 타입에 권장됩니다

모범 사례

  1. JSON 전용 데이터에는 JSON 타입 사용: JSON 객체만 저장한다면 기본 JSON 타입을 사용합니다(variant_types 속성 없음)
  2. 타입을 명시적으로 지정: Variant()를 사용할 때는 저장할 예정인 모든 타입을 명시적으로 나열합니다
  3. 실험적 기능 활성화: ClickHouse에서 allow_experimental_json_type = 1이 활성화되어 있는지 확인합니다
  4. 쓰기에는 JSON 포맷 사용: 더 나은 호환성을 위해 VariantType 데이터 쓰기에는 JSON 포맷을 권장합니다
  5. 쿼리 패턴 고려: JSON/Variant 타입은 효율적인 필터링을 위해 ClickHouse의 JSON 경로 쿼리를 지원합니다
  6. 성능을 위한 컬럼 힌트: ClickHouse에서 JSON 필드를 사용할 때 컬럼 힌트를 추가하면 쿼리 성능이 향상됩니다. 현재는 Spark를 통해 컬럼 힌트를 추가하는 기능을 지원하지 않습니다. 이 기능의 진행 상황은 GitHub issue #497에서 확인하십시오.

예시: 전체 워크플로

import org.apache.spark.sql.functions._

// ClickHouse에서 실험적 JSON 타입 활성화
spark.sql("SET allow_experimental_json_type = 1")

// Variant 컬럼이 있는 테이블 생성
spark.sql("""
  CREATE TABLE clickhouse.default.events (
    event_id BIGINT,
    event_time TIMESTAMP,
    event_data VARIANT
  ) USING clickhouse
  TBLPROPERTIES (
    'clickhouse.column.event_data.variant_types' = 'String, Int64, Bool, JSON',
    'engine' = 'MergeTree()',
    'order_by' = 'event_time'
  )
""")

// 혼합 타입 데이터 준비
val events = Seq(
  (1L, "2024-01-01 10:00:00", """{"action": "login", "user_id": 123}"""),
  (2L, "2024-01-01 10:05:00", """{"action": "purchase", "amount": 99.99}"""),
  (3L, "2024-01-01 10:10:00", """{"action": "logout", "duration": 600}""")
).toDF("event_id", "event_time", "json_data")

// VariantType으로 변환 후 저장
val variantEvents = events.select(
  col("event_id"),
  to_timestamp(col("event_time")).as("event_time"),
  parse_json(col("json_data")).as("event_data")
)

variantEvents.writeTo("clickhouse.default.events").append()

// 읽기 및 쿼리
val result = spark.sql("""
  SELECT event_id, event_time, event_data
  FROM clickhouse.default.events
  WHERE event_time >= '2024-01-01'
  ORDER BY event_time
""")

result.show(false)

구성

다음은 커넥터에서 조정할 수 있는 구성입니다.


기본값 설명 도입 버전
spark.clickhouse.ignoreUnsupportedTransform true ClickHouse는 cityHash64(col_1, col_2)와 같은 복잡한 표현식을 세그먼트 분할 키나 파티션 값으로 사용하는 것을 지원하지만, 현재 Spark는 이를 지원하지 않습니다. true로 설정하면 지원되지 않는 표현식을 무시하고 경고를 기록하며, 그렇지 않으면 예외를 발생시키고 즉시 실패합니다. 경고: spark.clickhouse.write.distributed.convertLocal=true인 경우 지원되지 않는 세그먼트 분할 키를 무시하면 데이터가 손상될 수 있습니다. 커넥터는 이를 검증하며 기본적으로 오류를 발생시킵니다. 이를 허용하려면 spark.clickhouse.write.distributed.convertLocal.allowUnsupportedSharding=true를 명시적으로 설정하십시오. 0.4.0
spark.clickhouse.read.compression.codec lz4 읽을 때 데이터 압축을 해제하는 데 사용하는 코덱입니다. 지원되는 코덱은 none, lz4입니다. 0.5.0
spark.clickhouse.read.distributed.convertLocal true 분산 테이블을 읽을 때 테이블 자체 대신 로컬 테이블을 읽습니다. true인 경우 spark.clickhouse.read.distributed.useClusterNodes는 무시됩니다. 0.1.0
spark.clickhouse.read.fixedStringAs binary ClickHouse FixedString 타입을 지정한 Spark 데이터 타입으로 읽습니다. 지원 타입: binary, string 0.8.0
spark.clickhouse.read.format json 읽기 시 사용할 직렬화 포맷입니다. 지원 포맷: json, binary 0.6.0
spark.clickhouse.read.runtimeFilter.enabled false 읽기 시 런타임 필터를 사용하도록 설정합니다. 0.8.0
spark.clickhouse.read.splitByPartitionId true true이면 파티션 값 대신 가상 컬럼 _partition_id를 사용해 입력 파티션 필터를 구성합니다. 파티션 값을 기준으로 SQL 프레디케이트를 조합할 때 알려진 문제가 있습니다. 이 기능을 사용하려면 ClickHouse 서버 v21.6+가 필요합니다. 0.4.0
spark.clickhouse.useNullableQuerySchema false true이면 테이블 생성 시 CREATE/REPLACE TABLE ... AS SELECT ...를 실행할 때 쿼리 스키마의 모든 필드를 널 허용으로 표시합니다. 참고로 이 구성에는 SPARK-43390(Spark 3.5에서 사용 가능)이 필요하며, 이 패치가 없으면 항상 true로 동작합니다. 0.8.0
spark.clickhouse.write.batchSize 10000 ClickHouse에 쓸 때 배치당 레코드 수입니다. 0.1.0
spark.clickhouse.write.compression.codec lz4 데이터를 쓸 때 압축에 사용하는 코덱입니다. 지원되는 코덱: none, lz4. 0.3.0
spark.clickhouse.write.distributed.convertLocal false 분산 테이블에 쓸 때는 해당 테이블 자체가 아니라 로컬 테이블에 씁니다. true로 설정하면 spark.clickhouse.write.distributed.useClusterNodes는 무시됩니다. 이렇게 하면 ClickHouse의 네이티브 라우팅을 우회하므로 Spark가 세그먼트 분할 키를 평가해야 합니다. 지원되지 않는 세그먼트 분할 표현식을 사용할 때는 데이터가 잘못 분산되는 문제를 조용히 넘기지 않도록 spark.clickhouse.ignoreUnsupportedTransformfalse로 설정하십시오. 0.1.0
spark.clickhouse.write.distributed.convertLocal.allowUnsupportedSharding false 세그먼트 분할 키가 지원되지 않는 경우 convertLocal=trueignoreUnsupportedTransform=true로 분산 테이블에 쓰는 것을 허용합니다. 이는 위험하며, 잘못된 세그먼트 분할로 인해 데이터가 손상될 수 있습니다. true로 설정하면 Spark가 지원되지 않는 세그먼트 분할 표현식을 평가할 수 없으므로, 쓰기 전에 데이터가 올바르게 정렬되고 세그먼트 분할되어 있는지 반드시 확인해야 합니다. 위험을 충분히 이해하고 데이터 분포를 검증한 경우에만 true로 설정하십시오. 기본적으로 이 조합은 오류를 발생시켜, 데이터가 조용히 손상되는 것을 방지합니다. 0.10.0
spark.clickhouse.write.distributed.useClusterNodes true 분산 테이블에 쓸 때 클러스터의 모든 노드에 기록합니다. 0.1.0
spark.clickhouse.write.format arrow 쓰기 시 사용할 직렬화 포맷입니다. 지원되는 포맷: json, arrow 0.4.0
spark.clickhouse.write.localSortByKey true true이면 쓰기 전에 정렬 키 기준으로 로컬 정렬을 수행합니다. 0.3.0
spark.clickhouse.write.localSortByPartition spark.clickhouse.write.repartitionByPartition true이면 쓰기 전에 파티션 기준으로 로컬 정렬을 수행합니다. 설정하지 않으면 spark.clickhouse.write.repartitionByPartition 값과 동일합니다. 0.3.0
spark.clickhouse.write.maxRetry 3 재시도 가능한 오류 코드로 인해 단일 배치 쓰기가 실패한 경우, 쓰기 작업을 재시도하는 최대 횟수입니다. 0.1.0
spark.clickhouse.write.repartitionByPartition true 쓰기 전에 ClickHouse 테이블의 분포에 맞도록 ClickHouse 파티션 키를 기준으로 데이터를 다시 파티셔닝할지 여부입니다. 0.3.0
spark.clickhouse.write.repartitionNum 0 쓰기 전에 ClickHouse 테이블의 분포에 맞게 데이터를 다시 파티셔닝해야 합니다. 이 설정으로 재파티셔닝 수를 지정하십시오. 값이 1보다 작으면 재파티셔닝이 필요하지 않음을 의미합니다. 0.1.0
spark.clickhouse.write.repartitionStrictly false true이면 Spark는 쓰기 시 레코드를 데이터 소스 테이블에 전달하기 전에, 필요한 분산 요건을 충족하도록 들어오는 레코드를 각 파티션에 엄격하게 분산합니다. 그렇지 않으면 Spark가 쿼리 속도를 높이기 위해 일부 최적화를 적용할 수 있지만, 이 경우 분산 요건이 깨질 수 있습니다. 참고로 이 구성은 SPARK-37523(Spark 3.4에서 사용 가능)가 필요하며, 이 패치가 없으면 항상 true로 동작합니다. 0.3.0
spark.clickhouse.write.retryInterval 10s 쓰기 재시도 간의 간격(초)입니다. 0.1.0
spark.clickhouse.write.retryableErrorCodes 241 쓰기 실패 시 ClickHouse 서버에서 반환되는 재시도 가능한 오류 코드입니다. 0.1.0

지원되는 데이터 타입

이 섹션에서는 Spark와 ClickHouse 간의 데이터 타입 매핑을 설명합니다. 아래 표는 ClickHouse에서 Spark로 데이터를 읽어올 때와 Spark에서 ClickHouse로 데이터를 삽입할 때의 데이터 타입 변환을 빠르게 참고할 수 있도록 정리한 것입니다.

ClickHouse에서 Spark로 데이터 읽기

ClickHouse 데이터 타입 Spark 데이터 타입 지원 여부 기본 타입 여부 비고
Nothing NullType
Bool BooleanType
UInt8, Int16 ShortType
Int8 ByteType
UInt16,Int32 IntegerType
UInt32,Int64, UInt64 LongType
Int128,UInt128, Int256, UInt256 DecimalType(38, 0)
Float32 FloatType
Float64 DoubleType
String, UUID, Enum8, Enum16, IPv4, IPv6 StringType
FixedString BinaryType, StringType 구성 READ_FIXED_STRING_AS로 제어됨
Decimal DecimalType precision과 scale은 Decimal128까지 지원
Decimal32 DecimalType(9, scale)
Decimal64 DecimalType(18, scale)
Decimal128 DecimalType(38, scale)
Date, Date32 DateType
DateTime, DateTime32, DateTime64 TimestampType
Array ArrayType 아니요 배열 요소 타입도 변환됨
Map MapType 아니요 키는 StringType으로 제한됨
IntervalYear YearMonthIntervalType(Year)
IntervalMonth YearMonthIntervalType(Month)
IntervalDay, IntervalHour, IntervalMinute, IntervalSecond DayTimeIntervalType 아니요 해당 인터벌 유형이 사용됨
JSON, Variant VariantType 아니요 Spark 4.0+ 및 ClickHouse 25.3+가 필요합니다. spark.clickhouse.read.jsonAs=string을 사용하면 StringType으로 읽을 수 있습니다
Object
Nested
Tuple StructType 아니요 이름이 지정된 Tuple과 이름이 없는 튜플을 모두 지원합니다. 이름이 지정된 Tuple은 필드 이름을 기준으로 struct 필드에 매핑되며, 이름이 없는 튜플은 _1, _2 등을 사용합니다. 중첩 struct와 널 허용 필드도 지원합니다
Point
Polygon
MultiPolygon
Ring
IntervalQuarter
IntervalWeek
Decimal256
AggregateFunction
SimpleAggregateFunction

Spark에서 ClickHouse로 데이터 삽입

Spark 데이터 타입 ClickHouse 데이터 타입 지원 여부 기본 타입 여부 참고
BooleanType Bool 0.9.0 버전부터 UInt8이 아닌 Bool 타입에 매핑됩니다
ByteType Int8
ShortType Int16
IntegerType Int32
LongType Int64
FloatType Float32
DoubleType Float64
StringType String
VarcharType String
CharType String
DecimalType Decimal(p, s) 정밀도와 소수 자릿수는 Decimal128까지 지원됩니다
DateType Date
TimestampType DateTime
ArrayType (list, tuple, or array) Array 아니요 배열의 요소 타입도 함께 변환됩니다
MapType Map 아니요 키는 StringType으로 제한됩니다
StructType Tuple 아니요 필드 이름이 있는 named Tuple로 변환됩니다.
VariantType JSON or Variant 아니요 Spark 4.0+ 및 ClickHouse 25.3+가 필요합니다. 기본값은 JSON 타입입니다. 여러 타입을 갖는 Variant를 지정하려면 clickhouse.column.<name>.variant_types 속성을 사용하십시오.
Object
Nested

기여 및 지원

프로젝트에 기여하거나 문제를 보고하려는 경우, 의견을 보내주시면 감사하겠습니다! 이슈를 등록하고, 개선 사항을 제안하고, pull request를 제출하려면 GitHub 리포지토리를 방문하십시오. 기여는 언제나 환영합니다! 시작하기 전에 리포지토리의 기여 가이드라인을 확인해 주십시오. ClickHouse Spark 커넥터 개선에 도움을 주셔서 감사합니다!

Navigation