이 커넥터는 고급 파티셔닝과 프레디케이트 푸시다운 같은 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 API와 TableProvider 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 릴리스 사용 지침을 따르십시오.
dependencies {
implementation("com.clickhouse.spark:clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}:{{ stable_version }}")
implementation("com.clickhouse:clickhouse-jdbc:{{ clickhouse_jdbc_version }}:all") { transitive = false }
}SNAPSHOT 버전을 사용하려면 Gradle에서 Sonatype의 SNAPSHOT 릴리스 사용 지침을 따르십시오.
libraryDependencies += "com.clickhouse" % "clickhouse-jdbc" % {{ clickhouse_jdbc_version }} classifier "all"
libraryDependencies += "com.clickhouse.spark" %% clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }} % {{ stable_version }}Spark의 셸 옵션(Spark SQL CLI, Spark Shell CLI, Spark Submit 명령)으로 작업할 때는 필요한 JAR 파일을 전달해 의존성을 등록할 수 있습니다.
$SPARK_HOME/bin/spark-sql \
--jars /path/clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}:{{ stable_version }}.jar,/path/clickhouse-jdbc-{{ clickhouse_jdbc_version }}-all.jarJAR 파일을 Spark 클라이언트 노드에 복사하지 않으려면 대신 다음을 사용할 수 있습니다.
--repositories https://{maven-central-mirror or private-nexus-repo} \
--packages com.clickhouse.spark:clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}:{{ stable_version }},com.clickhouse:clickhouse-jdbc:{{ clickhouse_jdbc_version }}참고: SQL 전용 사용 사례의 프로덕션 환경에는 Apache Kyuubi를 권장합니다.
라이브러리 다운로드
바이너리 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-sqlCLI 명령에 전달). - 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()val 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()Dataset<Row> 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()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()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()// 명시적으로 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()// 명시적으로 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()// Overwrite 모드(먼저 테이블을 TRUNCATE함)
df.write
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "my_table")
.mode("overwrite")
.save()// 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()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()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();
}object NativeSparkRead extends App {
val 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
val df = spark.sql("select * from clickhouse.default.example_table")
df.show()
spark.stop()
}from pyspark.sql import SparkSession
packages = [
"com.clickhouse.spark:clickhouse-spark-runtime-3.4_2.12:0.8.0",
"com.clickhouse:clickhouse-client:0.7.0",
"com.clickhouse:clickhouse-http-client:0.7.0",
"org.apache.httpcomponents.client5:httpclient5:5.2.1"
]
spark = (SparkSession.builder
.config("spark.jars.packages", ",".join(packages))
.getOrCreate())
spark.conf.set("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
spark.conf.set("spark.sql.catalog.clickhouse.host", "127.0.0.1")
spark.conf.set("spark.sql.catalog.clickhouse.protocol", "http")
spark.conf.set("spark.sql.catalog.clickhouse.http_port", "8123")
spark.conf.set("spark.sql.catalog.clickhouse.user", "default")
spark.conf.set("spark.sql.catalog.clickhouse.password", "123456")
spark.conf.set("spark.sql.catalog.clickhouse.database", "default")
spark.conf.set("spark.clickhouse.write.format", "json")
df = spark.sql("select * from clickhouse.default.example_table")
df.show() CREATE TEMPORARY VIEW jdbcTable
USING org.apache.spark.sql.jdbc
OPTIONS (
url "jdbc:ch://localhost:8123/default",
dbtable "schema.tablename",
user "username",
password "password",
driver "com.clickhouse.jdbc.ClickHouseDriver"
);
SELECT * FROM jdbcTable;데이터 쓰기
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();
}object NativeSparkWrite extends App {
// Spark 세션을 생성합니다.
val spark: SparkSession = 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의 스키마를 정의합니다.
val rows = Seq(Row(1, "John"), Row(2, "Doe"))
val schema = List(
StructField("id", DataTypes.IntegerType, nullable = false),
StructField("name", StringType, nullable = true)
)
// df를 생성합니다.
val df: DataFrame = spark.createDataFrame(
spark.sparkContext.parallelize(rows),
StructType(schema)
)
df.writeTo("clickhouse.default.example_table").append()
spark.stop()
}from pyspark.sql import SparkSession
from pyspark.sql import Row
# 위에 제공된 호환성 매트릭스를 충족하는 다른 패키지 조합을 사용해도 됩니다.
packages = [
"com.clickhouse.spark:clickhouse-spark-runtime-3.4_2.12:0.8.0",
"com.clickhouse:clickhouse-client:0.7.0",
"com.clickhouse:clickhouse-http-client:0.7.0",
"org.apache.httpcomponents.client5:httpclient5:5.2.1"
]
spark = (SparkSession.builder
.config("spark.jars.packages", ",".join(packages))
.getOrCreate())
spark.conf.set("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
spark.conf.set("spark.sql.catalog.clickhouse.host", "127.0.0.1")
spark.conf.set("spark.sql.catalog.clickhouse.protocol", "http")
spark.conf.set("spark.sql.catalog.clickhouse.http_port", "8123")
spark.conf.set("spark.sql.catalog.clickhouse.user", "default")
spark.conf.set("spark.sql.catalog.clickhouse.password", "123456")
spark.conf.set("spark.sql.catalog.clickhouse.database", "default")
spark.conf.set("spark.clickhouse.write.format", "json")
# DataFrame을 생성합니다.
data = [Row(id=11, name="John"), Row(id=12, name="Doe")]
df = spark.createDataFrame(data)
# DataFrame을 ClickHouse에 기록합니다.
df.writeTo("clickhouse.default.example_table").append() -- resultTable은 clickhouse.default.example_table에 삽입할 Spark 중간 DataFrame입니다.
INSERT INTO TABLE clickhouse.default.example_table
SELECT * FROM resultTable;
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의 JSON 및 Variant 타입에 매핑되므로, 유연한 스키마의 데이터를 효율적으로 저장하고 쿼리할 수 있습니다.
ClickHouse 타입 매핑
| ClickHouse 유형 | Spark 유형 | 설명 |
|---|---|---|
JSON |
VariantType |
JSON 객체만 저장합니다({로 시작해야 합니다) |
Variant(T1, T2, ...) |
VariantType |
기본 타입, 배열, JSON을 포함한 여러 타입을 저장합니다 |
VariantType 데이터 읽기
ClickHouse에서 데이터를 읽으면 JSON 및 Variant 컬럼이 자동으로 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()# JSON 컬럼을 VariantType으로 읽기
df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
# variant 데이터에 액세스
df.show()
# 확인할 수 있도록 variant를 JSON 문자열로 변환
from pyspark.sql.functions import to_json
df.select(
"id",
to_json("data").alias("data_json")
).show()// JSON 컬럼을 VariantType으로 읽기
Dataset<Row> df = spark.sql("SELECT id, data FROM clickhouse.default.json_table");
// variant 데이터에 액세스
df.show();
// 확인할 수 있도록 variant를 JSON 문자열로 변환
import static 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'
)
""")from pyspark.sql.functions import parse_json
# JSON 데이터로 DataFrame 생성
json_data = [
(1, '{"name": "Alice", "age": 30}'),
(2, '{"name": "Bob", "age": 25}'),
(3, '{"name": "Charlie", "city": "NYC"}')
]
df = spark.createDataFrame(json_data, ["id", "json_string"])
# JSON 문자열을 VariantType으로 파싱
variant_df = df.select(
"id",
parse_json("json_string").alias("data")
)
# JSON 타입으로 ClickHouse에 쓰기
variant_df.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'
)
""")import static org.apache.spark.sql.functions.*;
// JSON 데이터로 DataFrame 생성
List<Row> jsonData = Arrays.asList(
RowFactory.create(1, "{\"name\": \"Alice\", \"age\": 30}"),
RowFactory.create(2, "{\"name\": \"Bob\", \"age\": 25}"),
RowFactory.create(3, "{\"name\": \"Charlie\", \"city\": \"NYC\"}")
);
StructType schema = new StructType(new StructField[]{
DataTypes.createStructField("id", DataTypes.IntegerType, false),
DataTypes.createStructField("json_string", DataTypes.StringType, false)
});
Dataset<Row> jsonDF = spark.createDataFrame(jsonData, schema);
// JSON 문자열을 VariantType으로 파싱
Dataset<Row> variantDF = jsonDF.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으로 읽힘# JSON/Variant를 VariantType 대신 문자열로 읽기
spark.conf.set("spark.clickhouse.read.jsonAs", "string")
df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
# data 컬럼은 JSON 문자열을 포함하는 StringType으로 읽힘// JSON/Variant를 VariantType 대신 문자열로 읽기
spark.conf().set("spark.clickhouse.read.jsonAs", "string");
Dataset<Row> df = spark.sql("SELECT id, data FROM clickhouse.default.json_table");
// data 컬럼은 JSON 문자열을 포함하는 StringType으로 읽힘쓰기 포맷 지원
VariantType의 쓰기 지원은 포맷에 따라 다릅니다:
| 포맷 | 지원 | 참고 사항 |
|---|---|---|
| JSON | ✅ 전체 지원 | JSON 및 Variant 타입을 모두 지원합니다. VariantType 데이터에는 JSON 사용을 권장합니다 |
| Arrow | ⚠️ 부분 지원 | ClickHouse JSON 타입으로 쓰기를 지원합니다. ClickHouse Variant 타입은 지원하지 않습니다. 전체 지원은 https://github.com/ClickHouse/ClickHouse/issues/92752 이슈가 해결되면 제공될 예정입니다 |
쓰기 포맷을 설정합니다:
spark.conf.set("spark.clickhouse.write.format", "json") // Variant 타입에 권장됩니다모범 사례
- JSON 전용 데이터에는 JSON 타입 사용: JSON 객체만 저장한다면 기본 JSON 타입을 사용합니다(
variant_types속성 없음) - 타입을 명시적으로 지정:
Variant()를 사용할 때는 저장할 예정인 모든 타입을 명시적으로 나열합니다 - 실험적 기능 활성화: ClickHouse에서
allow_experimental_json_type = 1이 활성화되어 있는지 확인합니다 - 쓰기에는 JSON 포맷 사용: 더 나은 호환성을 위해 VariantType 데이터 쓰기에는 JSON 포맷을 권장합니다
- 쿼리 패턴 고려: JSON/Variant 타입은 효율적인 필터링을 위해 ClickHouse의 JSON 경로 쿼리를 지원합니다
- 성능을 위한 컬럼 힌트: 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)from pyspark.sql.functions import parse_json, to_timestamp
# 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'
)
""")
# 혼합 타입 데이터 준비
events = [
(1, "2024-01-01 10:00:00", '{"action": "login", "user_id": 123}'),
(2, "2024-01-01 10:05:00", '{"action": "purchase", "amount": 99.99}'),
(3, "2024-01-01 10:10:00", '{"action": "logout", "duration": 600}')
]
df = spark.createDataFrame(events, ["event_id", "event_time", "json_data"])
# VariantType으로 변환 후 저장
variant_events = df.select(
"event_id",
to_timestamp("event_time").alias("event_time"),
parse_json("json_data").alias("event_data")
)
variant_events.writeTo("clickhouse.default.events").append()
# 읽기 및 쿼리
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(truncate=False)import static 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'" +
")");
// 혼합 타입 데이터 준비
List<Row> events = Arrays.asList(
RowFactory.create(1L, "2024-01-01 10:00:00", "{\"action\": \"login\", \"user_id\": 123}"),
RowFactory.create(2L, "2024-01-01 10:05:00", "{\"action\": \"purchase\", \"amount\": 99.99}"),
RowFactory.create(3L, "2024-01-01 10:10:00", "{\"action\": \"logout\", \"duration\": 600}")
);
StructType eventSchema = new StructType(new StructField[]{
DataTypes.createStructField("event_id", DataTypes.LongType, false),
DataTypes.createStructField("event_time", DataTypes.StringType, false),
DataTypes.createStructField("json_data", DataTypes.StringType, false)
});
Dataset<Row> eventsDF = spark.createDataFrame(events, eventSchema);
// VariantType으로 변환 후 저장
Dataset<Row> variantEvents = eventsDF.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();
// 읽기 및 쿼리
Dataset<Row> 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.ignoreUnsupportedTransform을 false로 설정하십시오. |
0.1.0 |
| spark.clickhouse.write.distributed.convertLocal.allowUnsupportedSharding | false | 세그먼트 분할 키가 지원되지 않는 경우 convertLocal=true 및 ignoreUnsupportedTransform=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 커넥터 개선에 도움을 주셔서 감사합니다!