Sangwon
06/02/2026, 1:39 AMRocky
06/02/2026, 1:40 AMSTRING column (like dt) is to use Expression Partitioning (specifically Automatic Partitioning).
Since you want the column to be a STRING but behave like a date for partitioning (including automatic creation and deletion of old data), you should use the str2date function in your partition expression.
Example: CREATE TABLE with String dt and Dynamic Behavior
sql
CREATE TABLE internal_catalog.example_db.my_table (
dt STRING NOT NULL COMMENT "Partition column in YYYY-MM-DD format",
user_id INT NOT NULL,
event_type INT,
revenue DECIMAL(10, 2)
)
DUPLICATE KEY(dt, user_id) -- Duplicate Key model
PARTITION BY str2date(dt, '%Y-%m-%d') -- Expression Partitioning (auto-creates partitions)
DISTRIBUTED BY HASH(user_id) BUCKETS AUTO -- Hash distribution with auto-bucketing
PROPERTIES (
"partition_live_number" = "30" -- Automatically keeps the last 30 partitions (Dynamic cleanup)
);
Key Components Explained:
1. Internal Catalog: By specifying internal_catalog.database.table (or just database.table if you are already in the internal catalog), the table is created in StarRocks' native storage.
2. String Field with str2date:
* The dt field remains a STRING in your schema.
* PARTITION BY str2date(dt, '%Y-%m-%d') tells StarRocks to treat the string as a date for partitioning purposes. This enables Automatic Partitioning, meaning StarRocks will create new partitions on the fly as data is loaded.
3. Duplicate Key Model: The DUPLICATE KEY(dt, user_id) clause explicitly defines this as a Duplicate Key table. This is ideal for log data where you want to keep all records.
4. DISTRIBUTED BY HASH: Data is distributed across buckets based on the user_id. In version 4.0+, using BUCKETS AUTO is recommended as StarRocks will manage the bucket count based on the data volume.
5. Dynamic Cleanup: The property "partition_live_number" = "30" replaces the legacy dynamic partition start property. It ensures that only the most recent 30 partitions are retained, automatically dropping older ones based on the date value.
Why use Expression Partitioning instead of Legacy Dynamic Partitioning?
In StarRocks 4.0.x, the Legacy Dynamic Partitioning (using dynamic_partition.enable) is primarily intended for `DATE`/`DATETIME` columns and requires manual configuration of prefixes and history offsets.
Expression Partitioning is more flexible for STRING types because:
* It creates partitions on demand when data arrives.
* It supports functions like str2date and substr to handle non-standard formats.
* It simplifies management by eliminating the need for FE background scheduling parameters for partition-specific properties like "creating tomorrow's partition in advance."
References
• table_design/data_distribution/expression_partitioning.md
• administration/management/FE_configuration.mdRocky
06/02/2026, 2:28 AMSTRING 타입의 컬럼(dt)에 레거시 방식의 Dynamic Partition(동적 파티션) 설정을 적용하려면, 해당 문자열을 날짜로 변환하여 인식할 수 있도록 str2date 함수를 함께 사용해야 합니다.
다음은 요청하신 모든 설정값(dynamic_partition 속성들)을 포함한 예시입니다.
예시: Duplicate Key 모델 및 동적 파티션 설정
sql
CREATE TABLE internal_catalog.example_db.my_table (
dt STRING NOT NULL COMMENT "YYYY-MM-DD 형식의 날짜 문자열",
user_id INT NOT NULL,
event_type INT,
revenue DECIMAL(10, 2)
)
DUPLICATE KEY(dt, user_id) -- Duplicate Key 모델 설정
PARTITION BY RANGE(str2date(dt, '%Y-%m-%d')) (
-- 동적 파티션이 시작되기 위해 최소 하나의 초기 파티션이 필요합니다.
PARTITION p20240101 VALUES LESS THAN ("2024-01-02")
)
DISTRIBUTED BY HASH(user_id) BUCKETS 10 -- Hash 분포 설정
PROPERTIES (
"dynamic_partition.enable" = "true", -- 동적 파티션 기능 활성화
"dynamic_partition.time_unit" = "DAY", -- 파티션 생성 단위 (DAY, WEEK, MONTH 등)
"dynamic_partition.start" = "-30", -- 현재 기준 30일 이전 파티션 삭제 (데이터 정리)
"dynamic_partition.end" = "3", -- 현재 기준 3일 후의 파티션까지 미리 생성
"dynamic_partition.prefix" = "p", -- 파티션 이름의 접두사 (예: p20240102)
"dynamic_partition.buckets" = "10", -- 동적으로 생성될 파티션당 버킷 수
"dynamic_partition.history_partition_num" = "0" -- 테이블 생성 시 과거 파티션을 일괄 생성할 개수
);
주요 설정값(Properties) 상세 설명
1. dynamic_partition.enable: "true"로 설정하여 백그라운드에서 파티션을 자동으로 관리하도록 합니다.
2. dynamic_partition.time_unit: 파티션을 나누는 시간 단위입니다. DAY, WEEK, MONTH, HOUR, `YEAR`를 지원합니다.
3. dynamic_partition.start: 보관할 과거 데이터의 범위입니다. `-30`으로 설정하면 30일이 지난 파티션은 자동으로 삭제(Drop)됩니다.
4. dynamic_partition.end: 미리 생성해둘 미래 파티션의 범위입니다. `3`으로 설정하면 오늘 기준으로 3일 뒤까지의 파티션을 미리 만들어 둡니다.
5. dynamic_partition.prefix: 생성되는 파티션 이름의 앞글자입니다.
6. dynamic_partition.buckets: 동적으로 생성되는 각 파티션에 적용될 버킷(태블릿)의 개수입니다.
7. dynamic_partition.history_partition_num: 테이블 생성 시점에 과거 파티션을 미리 만들어두고 싶을 때 사용합니다.
참고 사항
* 초기 파티션 필수: 레거시 동적 파티션 방식은 CREATE TABLE 시점에 PARTITION BY RANGE 구문에 최소 한 개의 수동 파티션이 정의되어 있어야만 작동합니다.
* 권장 사항: StarRocks 4.0.4 버전에서는 위와 같은 복잡한 속성 설정 없이도 자동으로 파티션을 생성/삭제해주는 Expression Partitioning (Automatic Partitioning) 방식을 더 권장합니다. (첫 번째 답변에서 안내드린 partition_live_number 속성 방식이 이에 해당합니다.)
References
• table_design/data_distribution/dynamic_partitioning.md
• sql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md
References
• table_design/data_distribution/dynamic_partitioning.md
• table_design/data_distribution/feature-support-data-distribution.mdSangwon
06/02/2026, 2:41 AMRocky
06/02/2026, 2:41 AMSELECT 문을 사용합니다.
sql
SELECT * FROM my_table WHERE dt = '2024-01-01';
2. 파티션 프루닝(Partition Pruning) 활용
파티션 테이블을 조회할 때 가장 중요한 것은 *성능*입니다. WHERE 절에 파티션 키(예: dt)를 조건으로 넣으면, StarRocks는 전체 데이터를 스캔하지 않고 해당 파티션의 데이터만 읽습니다. 이를 *파티션 프루닝*이라고 합니다.
* 효율적인 조회: 파티션 컬럼을 조건에 포함하세요.
sql
-- 특정 날짜 범위 조회 (해당 파티션만 접근)
SELECT user_id, SUM(revenue)
FROM my_table
WHERE dt >= '2024-01-01' AND dt <= '2024-01-07'
GROUP BY user_id;
3. 파티션 목록 및 정보 조회 (관리용)
현재 테이블에 어떤 파티션들이 생성되어 있고, 데이터가 얼마나 들어있는지 확인하고 싶을 때 사용합니다.
* 기본 조회: 생성된 파티션 이름, 범위(Range), 데이터 크기 등을 확인합니다.
sql
SHOW PARTITIONS FROM my_table;
* 시스템 뷰 조회 (v3.1.9+ / v3.2.4+): SQL을 통해 더 복잡한 조건으로 파티션을 조회할 수 있습니다.
sql
SELECT * FROM information_schema.partitions_meta
WHERE TABLE_NAME = 'my_table'
ORDER BY DATA_SIZE DESC;
4. 쿼리가 파티션을 잘 타는지 확인 (EXPLAIN)
내가 작성한 쿼리가 실제로 파티션 프루닝을 수행하는지 확인하려면 쿼리 앞에 `EXPLAIN`을 붙여 실행 계획을 확인하세요.
sql
EXPLAIN SELECT * FROM my_table WHERE dt = '2024-01-01';
* 결과 내용 중 OlapScanNode 부분의 partitions 항목을 확인하세요.
* 예: partitions=1/30 (전체 30개 파티션 중 1개만 읽음 → 정상 동작)
요약
* 데이터 조회: SELECT ... WHERE 파티션키 = '값'
* 파티션 상태 확인: SHOW PARTITIONS FROM 테이블명
* 성능 확인: EXPLAIN 실행 후 partitions 수치 확인
References
• sql-reference/sql-statements/table_bucket_part_index/SHOW_PARTITIONS.md
• best_practices/partitioning.md
• table_design/data_distribution/expression_partitioning.md
References
• sql-reference/sql-statements/table_bucket_part_index/SHOW_PARTITIONS.md
• knowledge/admin_and_operation/shared_data_compaction.mdSangwon
06/02/2026, 2:47 AMRocky
06/02/2026, 2:47 AMdt, layer_nm, fab_id 세 개의 컬럼을 조합하여 파티션을 구성하는 것이 가능합니다.
가장 현대적이고 권장되는 방식인 *자동 파티셔닝(Automatic Partitioning)*을 사용하여 예시를 구성해 드립니다. 이 방식을 사용하면 데이터가 적재될 때 (날짜 + 레이어 + Fab) 조합에 맞는 파티션이 자동으로 생성됩니다.
여러 컬럼을 사용하는 파티션 테이블 예시
sql
CREATE TABLE internal_catalog.example_db.multi_partition_table (
dt STRING NOT NULL COMMENT "YYYY-MM-DD 형식의 날짜",
layer_nm STRING NOT NULL,
fab_id STRING NOT NULL,
user_id INT NOT NULL,
revenue DECIMAL(10, 2)
)
DUPLICATE KEY(dt, layer_nm, fab_id, user_id) -- Duplicate Key 모델
-- str2date 함수와 나머지 컬럼들을 조합하여 파티션 설정
PARTITION BY str2date(dt, '%Y-%m-%d'), layer_nm, fab_id
DISTRIBUTED BY HASH(user_id) BUCKETS AUTO
PROPERTIES (
-- 주의: 여러 컬럼 조합 시 partition_live_number는 전체 파티션 개수를 기준으로 작동합니다.
-- 만약 30일치를 유지하고 싶고 layer_nm이 2개, fab_id가 2개라면 총 30*2*2 = 120개를 설정해야 합니다.
"partition_live_number" = "120"
);
상세 설명 및 주의사항
1. 복합 파티션 구성:
* PARTITION BY str2date(dt, '%Y-%m-%d'), layer_nm, fab_id 구문을 통해 세 가지 조건이 결합된 파티션이 만들어집니다.
* 예를 들어, dt='2024-01-01', layer_nm='PHOTO', fab_id='FAB1' 데이터가 들어오면 p20240101_PHOTO_FAB1 형태의 파티션이 자동으로 생성됩니다.
2. 컬럼 순서:
* `PARTITION BY`에 나열된 순서대로 파티션 이름과 구조가 결정됩니다. 가장 범위가 큰 단위(보통 날짜)를 앞에 두는 것이 일반적입니다.
3. 데이터 타입:
* `dt`는 `STRING`이지만 str2date 함수를 통해 DATE 타입으로 변환되어 파티션 키로 사용됩니다.
* `layer_nm`과 `fab_id`는 문자열 그대로 파티션 키로 활용됩니다.
4. 파티션 개수 관리 (partition_live_number):
* 단일 날짜 파티션과 달리, 여러 컬럼을 조합하면 생성되는 파티션의 총 개수가 급격히 늘어날 수 있습니다(날짜 수 × 레이어 수 × Fab 수).
* StarRocks는 너무 많은 파티션이 생기는 것을 방지하기 위해 기본적으로 한 번의 로드 작업당 생성 가능한 파티션 수를 제한합니다(기본 4096개).
* TTL(자동 삭제): partition_live_number 속성을 사용하면 지정된 개수보다 오래된 파티션을 자동으로 삭제합니다. 이때 기준은 "파티션 생성 시간"이 아닌 파티션 키의 값(날짜 등)을 기준으로 가장 최신 것들을 남깁니다.
쿼리 성능 최적화 (Partition Pruning)
조회 시 WHERE 절에 세 가지 조건을 모두 넣으면 성능이 극대화됩니다.
sql
SELECT * FROM multi_partition_table
WHERE dt = '2024-01-01'
AND layer_nm = 'PHOTO'
AND fab_id = 'FAB1';
이렇게 조회하면 StarRocks는 수많은 파티션 중 정확히 해당 조합의 파티션 하나만 읽으므로 매우 빠릅니다.
References
• table_design/data_distribution/expression_partitioning.md
• table_design/data_distribution/feature-support-data-distribution.md
• sql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md
References
• table_design/data_distribution/feature-support-data-distribution.md
• table_design/data_distribution/Data_distribution.md
• table_design/data_distribution/list_partitioning.md
• table_design/data_distribution/dynamic_partitioning.mdSangwon
06/02/2026, 2:49 AMRocky
06/02/2026, 2:49 AMsql
CREATE TABLE internal_catalog.example_db.permanent_multi_partition_table (
dt STRING NOT NULL,
layer_nm STRING NOT NULL,
fab_id STRING NOT NULL,
user_id INT NOT NULL,
revenue DECIMAL(10, 2)
)
DUPLICATE KEY(dt, layer_nm, fab_id, user_id)
-- 자동 파티션 설정: dt, layer_nm, fab_id의 조합으로 파티션 생성
PARTITION BY str2date(dt, '%Y-%m-%d'), layer_nm, fab_id
DISTRIBUTED BY HASH(user_id) BUCKETS AUTO
PROPERTIES (
-- "partition_live_number" 속성을 아예 생략합니다.
-- 이 경우, 생성된 파티션은 자동으로 삭제되지 않고 계속 누적됩니다.
"replicated_storage" = "true"
);
주요 변경 사항 및 설명
1. partition_live_number 삭제:
* 이 속성이 없으면 StarRocks는 파티션 개수를 관리(삭제)하지 않습니다.
* 데이터가 들어오는 대로 (날짜, 레이어, Fab) 조합의 파티션이 계속해서 늘어나며, 과거 데이터는 그대로 보존됩니다.
2. 레거시 방식(dynamic_partition)과의 차이:
* 만약 레거시 방식인 dynamic_partition.enable = "true" 설정을 사용 중이라면, dynamic_partition.start 값을 설정하지 않거나 매우 큰 음수 값(예: "-99999")으로 설정해야 삭제를 피할 수 있습니다.
* 하지만 자동 파티셔닝(Expression Partitioning) 방식에서는 위 예시처럼 관련 속성을 비워두는 것이 가장 깔끔하고 권장되는 방법입니다.
3. 관리 시 주의사항:
* 파티션 개수 모니터링: 데이터를 영구히 보관하면 파티션 수가 수만 개 이상으로 늘어날 수 있습니다. StarRocks는 많은 수의 파티션을 잘 처리하지만, 너무 잘게 쪼개질 경우(예: 시간 단위 + 여러 컬럼 조합) 메타데이터 관리에 부담이 될 수 있으므로 주기적으로 `SHOW PARTITIONS FROM 테이블명`을 통해 파티션 개수를 확인해 보시는 것이 좋습니다.
* 수동 삭제: 만약 나중에 특정 시점의 데이터가 정말 필요 없어져서 지우고 싶다면, 아래 명령어로 특정 파티션만 골라서 수동으로 삭제할 수 있습니다.
sql
ALTER TABLE permanent_multi_partition_table DROP PARTITION <파티션이름>;
결론적으로, `partition_live_number`나 dynamic_partition.start 같은 설정만 넣지 않으시면 데이터는 삭제되지 않고 계속 쌓이 유지됩니다.
References
• table_design/data_distribution/expression_partitioning.md
• sql-reference/sql-statements/table_bucket_part_index/ALTER_TABLE.md
• table_design/data_distribution/dynamic_partitioning.mdSangwon
06/02/2026, 5:00 AMRocky
06/02/2026, 5:00 AMscala
df.write
.format("starrocks")
.option("starrocks.fe.http.url", "127.0.0.1:8030")
.option("starrocks.fe.jdbc.url", "jdbc:<mysql://127.0.0.1:9030%22|mysql://127.0.0.1:9030">) // 메타데이터 조회용
.option("starrocks.table.identifier", "db_name.table_name")
.option("starrocks.user", "root")
.option("starrocks.password", "")
.mode("append")
.save()
2. 특정 스키마 강제 지정하기 (starrocks.column.types)
만약 Spark가 추론한 타입과 StarRocks의 타입이 맞지 않거나(예: BITMAP, HLL 타입 등), 특정 컬럼만 명시하고 싶을 때 사용합니다.
* 설정값: col_name1 INT, col_name2 STRING 형식의 DDL 포맷
* 사용 예시:
scala
.option("starrocks.column.types", "user_id INT, dt DATE, revenue DECIMAL(10,2)")
3. 프로그램 코드 내에서 스키마 확인하기
현재 StarRocks 테이블의 스키마가 Spark에서 어떻게 인식되는지 미리 확인하고 싶다면, `load()`를 통해 빈 데이터를 읽어와서 출력해 볼 수 있습니다.
scala
val srTable = spark.read
.format("starrocks")
.option("starrocks.fe.http.url", "127.0.0.1:8030")
.option("starrocks.fe.jdbc.url", "jdbc:<mysql://127.0.0.1:9030%22|mysql://127.0.0.1:9030">)
.option("starrocks.table.identifier", "db_name.table_name")
.option("starrocks.user", "root")
.option("starrocks.password", "")
.load()
srTable.printSchema() // StarRocks 테이블의 스키마 출력
주의사항
* 컬럼명 매핑: Spark DataFrame의 컬럼명과 StarRocks 테이블의 컬럼명이 대소문자까지 일치해야 합니다. (일치하지 않을 경우 starrocks.columns 옵션으로 수동 매핑 가능)
* JDBC 드라이버: v1.1.1 버전부터는 mysql-connector-java 라이브러리가 커넥터에 포함되지 않으므로, Spark 실행 시 해당 jar 파일을 반드시 클래스패스에 추가해야 합니다.
References
• loading/Spark-connector-starrocks.md
• unloading/Spark_connector.md
References
• loading/Spark-connector-starrocks.md
• unloading/Spark_connector.md
• loading/Flink-connector-starrocks.md
• unloading/Flink_connector.md
References
• loading/Spark-connector-starrocks.md
• unloading/Spark_connector.mdSangwon
06/02/2026, 5:04 AMRocky
06/02/2026, 5:04 AMload() 기능을 사용하면 실제 데이터를 읽지 않고도 스키마 정보만 메타데이터 형태로 가져올 수 있습니다.
[PySpark 예시]
python
# 1. StarRocks 테이블의 스키마 정보만 가져오기 (데이터는 읽지 않음)
target_table_df = spark.read \
.format("starrocks") \
.option("starrocks.fe.http.url", "127.0.0.1:8030") \
.option("starrocks.fe.jdbc.url", "jdbc:<mysql://127.0.0.1:9030%22|mysql://127.0.0.1:9030">) \
.option("starrocks.table.identifier", "db_name.table_name") \
.option("starrocks.user", "root") \
.option("starrocks.password", "password") \
.load()
# 2. StarRocks 테이블의 컬럼명 리스트 추출
sr_columns = target_table_df.columns
# 3. 내 DataFrame(source_df)의 컬럼을 StarRocks 순서와 이름에 맞게 선택/변경
# (만약 컬럼명이 이미 같다면 순서만 맞추기)
final_df = source_df.select(*sr_columns)
# 4. 저장
final_df.write.format("starrocks")...save()
2. 컬럼명이 다를 경우 매핑하기 (Renaming)
만약 내 DataFrame의 컬럼명과 StarRocks의 컬럼명이 완전히 다르다면, `withColumnRenamed`나 `alias`를 사용해야 합니다.
python
# 예: Spark에는 'user_id'가 있는데 StarRocks에는 'id'인 경우
final_df = source_df.select(
col("user_id").alias("id"),
col("dt"),
col("layer_nm")
)
3. 커넥터 옵션 활용하기 (starrocks.columns)
DataFrame의 컬럼명을 직접 바꾸는 대신, StarRocks로 보낼 때 *어떤 컬럼을 어떤 순서로 넣을지 명시*할 수도 있습니다. 이 옵션을 쓰면 Spark DataFrame의 순서가 StarRocks와 달라도 매핑이 가능합니다.
python
.option("starrocks.columns", "id, dt, layer_nm, fab_id")
# Spark DataFrame의 첫 번째 컬럼이 StarRocks의 'id'로, 두 번째가 'dt'로 들어갑니다.
추천하는 작업 흐름
1. 순서 맞추기: StarRocks는 기본적으로 DataFrame의 컬럼 순서대로 데이터를 밀어 넣습니다. 따라서 `target_table_df.columns`를 사용하여 순서를 맞추는 것이 가장 안전합니다.
2. 대소문자 확인: StarRocks는 컬럼명 대소문자를 구분할 수 있으므로, 반드시 일치시켜야 합니다.
3. 불필요한 컬럼 제거: Spark DataFrame에 StarRocks 테이블에 없는 컬럼이 있으면 에러가 발생하므로, `select`를 통해 필요한 컬럼만 남기세요.
요약하자면, `spark.read...load().columns`로 StarRocks의 컬럼 리스트를 얻은 다음, `source_df.select(*cols)`를 적용하는 것이 가장 자동화된 방법입니다.
References
• loading/Spark-connector-starrocks.md
• unloading/Spark_connector.md
• sql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.mdSangwon
06/02/2026, 5:07 AMRocky
06/02/2026, 5:07 AMwrite 하기 전에 Spark의 cast() 함수를 사용하여 StarRocks 테이블의 데이터 타입과 일치시키는 방법입니다. Spark 수준에서 데이터를 검증할 수 있어 가장 안전합니다.
python
from pyspark.sql.functions import col
# 예: 'price' 컬럼을 Decimal로, 'created_at'을 Date 타입으로 변환
transformed_df = source_df.withColumn("price", col("price").cast("decimal(10,2)")) \
.withColumn("created_at", col("created_at").cast("date"))
transformed_df.write \
.format("starrocks") \
.option("starrocks.table.identifier", "db.table") \
...
.save()
2. starrocks.columns 옵션에서 변환 식 사용 (Stream Load 방식)
StarRocks Spark 커넥터는 내부적으로 *Stream Load*를 사용합니다. 따라서 starrocks.columns 옵션에 StarRocks의 함수를 직접 작성하여 변환할 수 있습니다.
* 용도: 소스 데이터가 String인데 StarRocks에 저장할 때 특정 함수(예: to_bitmap, str2date)를 거쳐야 하는 경우 유용합니다.
* 예시:
python
.option("starrocks.columns", "user_id, dt_str, date=str2date(dt_str, '%Y%m%d')")
_위 예시는 `dt_str`이라는 임시 컬럼으로 데이터를 받아 StarRocks의 date 컬럼에 변환해서 넣는 방식입니다._
3. starrocks.column.types 명시 (Schema Override)
v1.1.1 버전부터 지원되는 옵션으로, 커넥터가 StarRocks 테이블로부터 자동으로 스키마를 유추하는 대신 사용자가 정의한 타입을 강제로 적용합니다. BITMAP이나 HLL 같은 특수 타입을 처리할 때 주로 사용합니다.
python
.option("starrocks.column.types", "id INT, revenue DECIMAL(10,2), tags STRING")
주의할 타입 매핑
캐스팅 시 아래의 Spark-StarRocks 기본 매핑 관계를 참고하세요.
Spark 데이터 타입 | StarRocks 데이터 타입 | 비고
StringType | VARCHAR / STRING / JSON |
TimestampType | DATETIME |
DateType | DATE |
DecimalType | DECIMAL |
LongType | BIGINT |
팁: 만약 데이터 타입이 맞지 않으면 Stream Load 과정에서 "Quality not good enough" 에러가 발생하거나 데이터가 `NULL`로 들어갈 수 있습니다. 따라서 가급적 Spark에서 `cast()`를 통해 타입을 명확히 맞춘 후 전송하는 것이 디버깅에 유리합니다.
References
• loading/Spark-connector-starrocks.md
• unloading/Spark_connector.md
• sql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md
References
• unloading/Spark_connector.md
• loading/Spark-connector-starrocks.md
• unloading/Flink_connector.mdSangwon
06/02/2026, 5:30 AMRocky
06/02/2026, 5:31 AMappend 모드로 쓰면 자동으로 Upsert가 수행됩니다. PK가 동일한 행이 들어오면 기존 데이터를 덮어쓰고, 없으면 새로 추가합니다.
python
df.write \
.format("starrocks") \
.option("starrocks.fe.http.url", "127.0.0.1:8030") \
.option("starrocks.fe.jdbc.url", "jdbc:<mysql://127.0.0.1:9030%22|mysql://127.0.0.1:9030">) \
.option("starrocks.table.identifier", "test_db.score_board") \
.option("starrocks.user", "root") \
.option("starrocks.password", "password") \
.mode("append") \
.save()
2. 부분 업데이트 (Partial Update)
테이블의 전체 컬럼이 아니라 *특정 컬럼만 업데이트*하고 싶을 때 사용합니다. 예를 들어, id, name, score 컬럼 중 `id`와 `name`만 업데이트하고 싶다면 다음과 같이 설정합니다.
* `starrocks.write.properties.partial_update`: `true`로 설정
* `starrocks.columns`: 업데이트할 컬럼들 명시 (PK 포함 필수)
python
df.select("id", "name").write \
.format("starrocks") \
.option("starrocks.write.properties.partial_update", "true") \
.option("starrocks.columns", "id,name") \
.option("starrocks.table.identifier", "test_db.score_board") \
...
.save()
3. 조건부 업데이트 (Conditional Update)
데이터가 들어올 때, 특정 컬럼(예: updated_at 타임스탬프)의 값이 기존 값보다 *클 때만 업데이트*하도록 제한할 수 있습니다. 데이터의 순서가 뒤섞여 들어오는 경우에 유용합니다.
* `starrocks.write.properties.merge_condition`: 기준이 될 컬럼명 지정
python
.option("starrocks.write.properties.merge_condition", "update_time")
필수 체크리스트
1. 테이블 모델 확인: 반드시 `ENGINE=OLAP PRIMARY KEY(pk_col)`로 생성된 테이블이어야 합니다. (Duplicate/Aggregate 모델은 Upsert 미지원)
2. PK 포함: 데이터프레임에는 반드시 테이블의 Primary Key 컬럼이 포함되어 있어야 합니다.
3. *Operation 지정 (*op)*: 특정 행을 삭제(delete)하거나 명시적으로 연산을 지정하고 싶다면 *op 컬럼을 활용할 수 있으나, 일반적인 Upsert는 설정 없이 `append`만으로 충분합니다.
References
• loading/Spark-connector-starrocks.md
• loading/Load_to_Primary_Key_tables.md
• table_design/table_types/primary_key_table.md
References
• loading/Spark-connector-starrocks.md
• loading/Load_to_Primary_Key_tables.md
• table_design/table_types/primary_key_table.md
• loading/Flink-connector-starrocks.md
• unloading/Spark_connector.md
References
• loading/Spark-connector-starrocks.md
• loading/Load_to_Primary_Key_tables.md
• table_design/table_types/primary_key_table.md
• loading/Flink-connector-starrocks.md
• unloading/Spark_connector.md
• unloading/Flink_connector.md
References
• loading/Spark-connector-starrocks.md
• loading/Kafka-connector-starrocks.mdSangwon
06/02/2026, 5:32 AMRocky
06/02/2026, 5:32 AMdt, layer, fab 세 가지 컬럼으로 파티셔닝을 하는 예제입니다.
주의사항
1. PK 제약 조건: Primary Key 모델에서 파티션 키(`PARTITION BY`에 사용되는 컬럼)는 반드시 *Primary Key 컬럼 구성 요소에 포함*되어야 합니다.
2. 자동 파티션 생성: dt, layer, fab 조합이 많을 경우 일일이 파티션을 지정하기 힘들기 때문에, List 자동 파티션(Automatic Partitioning) 기능을 사용하는 것이 효율적입니다.
테이블 생성 SQL 예제
sql
CREATE TABLE my_db.my_table (
dt DATE NOT NULL,
layer VARCHAR(50) NOT NULL,
fab VARCHAR(50) NOT NULL,
id BIGINT NOT NULL,
value DOUBLE,
updated_at DATETIME
) ENGINE=OLAP
PRIMARY KEY(dt, layer, fab, id) -- 파티션 키(dt, layer, fab)가 PK에 포함되어야 함
PARTITION BY (dt, layer, fab) -- dt, layer, fab 기준으로 파티션 생성
DISTRIBUTED BY HASH(id) BUCKETS 10
PROPERTIES (
"replication_num" = "3",
"enable_persistent_index" = "true" -- PK 테이블 성능 최적화
);
주요 포인트 설명
* PARTITION BY (dt, layer, fab): 위와 같이 별도의 (PARTITION p1 VALUES IN (...)) 선언 없이 컬럼만 명시하면, 데이터가 삽입될 때 새로운 (dt, layer, fab) 조합이 들어오면 StarRocks가 *자동으로 파티션을 생성*합니다. (버전 3.1 이상 권장)
* PRIMARY KEY(dt, layer, fab, id): dt, layer, `fab`를 파티션 키로 썼으므로 이들은 자동으로 PK의 일부가 되어야 합니다. `id`는 행을 식별하기 위한 추가 키입니다.
* DISTRIBUTED BY HASH(id): 데이터 분산(Bucketing)은 고유값이 많은 `id`나 다른 PK 컬럼을 사용하는 것이 데이터 쏠림 방지에 좋습니다.
만약 버전이 낮아 자동 파티션이 안 될 경우 (수동 List 파티션)
특정 값들을 미리 지정해야 한다면 아래와 같은 형식을 사용합니다.
sql
PARTITION BY LIST (dt, layer, fab) (
PARTITION p20231001_layer1_fabA VALUES IN (("2023-10-01", "L1", "FAB_A")),
PARTITION p20231001_layer1_fabB VALUES IN (("2023-10-01", "L1", "FAB_B"))
)
*일반적으로 dt, layer, fab 조합은 가변적이므로 첫 번째 예제인 자동 파티션 방식을 가장 추천합니다.*
References
• sql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md
• table_design/table_types/primary_key_table.md
• table_design/partitioning/automatic_partitioning.mdSangwon
06/02/2026, 5:33 AMRocky
06/02/2026, 5:33 AMdt, layer, fab 예제를 다시 보면:
sql
PRIMARY KEY(dt, layer, fab, id) -- 이 4개 조합이 PK
PARTITION BY (dt, layer, fab) -- PK 구성 요소 중 일부이므로 '가능'
만약 아래처럼 선언하면 *에러*가 발생합니다.
sql
PRIMARY KEY(id) -- PK는 id 하나뿐인데
PARTITION BY (dt) -- PK에 없는 dt로 파티션을 나누려 하면 '에러'
따라서 PK 테이블에서 파티션을 나누고 싶다면, *해당 파티션 컬럼들을 모두 PK 선언부(괄호 안)에 포함*시켜 주어야 합니다.
References
• table_design/partitioning/table_partitioning.md
• table_design/table_types/primary_key_table.md
• sql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.mdSangwon
06/02/2026, 5:34 AMRocky
06/02/2026, 5:34 AMappend 모드로 데이터를 넣으면, 기존에 동일한 ID나 키가 있더라도 상관하지 않고 *매번 새로운 행으로 추가*합니다.
원하시는 Upsert(있으면 Update, 없으면 Insert) 기능을 사용하시려면 반드시 *Primary Key 모델*을 사용해야 합니다.
왜 Duplicate Key에서는 안 되나요?
* Duplicate Key 모델: 데이터가 들어오는 대로 단순히 쌓는 구조입니다. 분석용 원천 로그 데이터를 그대로 저장할 때 사용하며, '업데이트'라는 개념 자체가 없습니다.
* Primary Key 모델: 데이터를 넣을 때 메모리에 있는 인덱스(Persistent Index)를 확인하여, 동일한 키가 있으면 해당 위치의 데이터를 수정하고 없으면 새로 넣습니다.
해결 방법: 테이블 모델 변경
테이블을 생성할 때 아래와 같이 PRIMARY KEY 모델로 만드셔야 합니다.
sql
CREATE TABLE my_db.my_upsert_table (
dt DATE NOT NULL,
layer VARCHAR(50) NOT NULL,
fab VARCHAR(50) NOT NULL,
id BIGINT NOT NULL,
value DOUBLE
) ENGINE=OLAP
PRIMARY KEY(dt, layer, fab, id) -- 여기에 지정된 컬럼들이 중복 체크의 기준이 됨
PARTITION BY (dt, layer, fab)
DISTRIBUTED BY HASH(id);
모델별 비교
기능 | Duplicate Key | Primary Key
중복 허용 | 허용 (모든 데이터 저장) | 불허 (동일 키는 최신 데이터로 덮어씀)
Upsert 지원 | 지원 안 함 (무조건 Insert) | 지원 함
주요 용도 | 단순 로그 저장, 원천 데이터 적재 | 실시간 대시보드, 수정이 빈번한 데이터
결론:
Spark에서 `append`를 썼을 때 Upsert로 동작하게 하려면, 대상 StarRocks 테이블이 *Primary Key*로 생성되어 있어야 합니다. 지금 사용 중인 Duplicate Key 테이블은 Upsert가 불가능한 구조이므로 테이블 정의를 변경하셔야 합니다.
References
• table_design/table_types/duplicate_key_table.md
• table_design/table_types/primary_key_table.md
• loading/Spark-connector-starrocks.mdSangwon
06/02/2026, 6:40 AMRocky
06/02/2026, 6:40 AMWHERE dt = '...')으로 조회할 때 검색 속도가 매우 빠릅니다. (유니크 제약과는 무관합니다.)
3. 가장 빠른 로딩 속도
* 새 데이터를 넣을 때 기존 데이터와 비교하거나 중복 체크를 하는 과정이 전혀 없습니다.
* 따라서 4가지 테이블 모델(Duplicate, Primary, Unique, Aggregate) 중 쓰기(Ingestion) 성능이 가장 뛰어납니다.
4. 언제 사용하나요?
* 원천 로그 데이터: 클릭 로그, 센서 데이터처럼 중복 삭제가 필요 없고 발생한 그대로를 저장해야 하는 경우.
* 단순 적재: 데이터 수정이 거의 없고, 오로지 빠르게 쌓고 조회(Filter)하는 것이 중요한 경우.
요약 예시
sql
CREATE TABLE logs (
dt DATE,
id INT,
event VARCHAR(20)
)
DUPLICATE KEY(dt, id) -- 이 컬럼들로 정렬만 수행함
...
위 테이블에 `(2023-10-01, 1, 'click')`을 두 번 INSERT하면, 테이블에는 똑같은 데이터가 2줄 남게 됩니다.
References
• table_design/table_types/duplicate_key_table.md
• table_design/table_types/table_types.mdSangwon
06/04/2026, 5:49 AMRocky
06/04/2026, 5:49 AMDUPLICATE KEY, PRIMARY KEY, UNIQUE KEY 등에 정의된 컬럼 순서대로 데이터가 정렬됩니다.
* 추천: WHERE 절에 가장 자주 쓰이는 컬럼, 변별력(Cardinality)이 높은 컬럼을 *가장 앞쪽*에 배치하세요.
2. Bloom Filter 인덱스 (bloom_filter_columns)
특정 값이 데이터에 포함되어 있는지 빠르게 확인하여 불필요한 파일을 읽지 않게 합니다.
* 용도: `id`나 serial_number 같이 값의 종류가 매우 많고(High Cardinality), = 또는 IN 조건으로 자주 조회하는 컬럼에 효과적입니다.
* 제약: TINYINT, FLOAT, DOUBLE, DECIMAL 타입은 지원하지 않습니다.
* 예시:
sql
PROPERTIES (
"bloom_filter_columns" = "user_id, order_id"
)
3. 데이터 압축 방식 (compression)
저장 공간을 아끼고 I/O를 줄여 성능을 높입니다.
* 추천:
* LZ4 (기본값): 압축/해제 속도가 매우 빨라 CPU 부하가 적습니다.
* `ZSTD`: 압축률이 매우 좋아 저장 공간을 크게 절약하며, I/O가 병목인 환경에서 성능이 더 좋을 수 있습니다.
* 예시:
sql
PROPERTIES (
"compression" = "ZSTD"
)
4. Inverted Index (역색인, v3.3+ 권장)
텍스트 검색이나 복합적인 필터링 성능을 획기적으로 높여줍니다.
* 용도: 로그 데이터의 텍스트 검색이나 여러 컬럼을 조합한 필터링이 많을 때 유리합니다.
* 예시:
sql
INDEX index_name (column_name) USING INVERTED [COMMENT 'comment']
5. Colocate Join (colocate_with)
자주 Join하는 두 테이블의 데이터를 물리적으로 동일한 노드에 배치하여, 조인 시 네트워크 전송(Shuffle)을 없앱니다.
* 설정 방법: Join 대상 테이블들에 동일한 그룹명을 지정합니다.
* 예시:
sql
PROPERTIES (
"colocate_with" = "group_name_1"
)
6. Bitmap Index
값의 종류가 적은(Low Cardinality) 컬럼에 대해 필터링 속도를 높입니다.
* 용도: gender, status, region 등 값의 종류가 몇 개 안 되는 컬럼에 적합합니다.
──────────
종합 추천 예시 SQL
sql
CREATE TABLE sales_records (
order_date DATE NOT NULL,
order_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
category VARCHAR(20),
amount DECIMAL(10, 2)
) ENGINE=OLAP
PRIMARY KEY(order_date, order_id) -- 1. 정렬 키 설정
PARTITION BY (order_date) -- 파티셔닝 (Pruning 효과)
DISTRIBUTED BY HASH(order_id) -- 버케팅
PROPERTIES (
"bloom_filter_columns" = "user_id", -- 2. 고변별력 컬럼 인덱스
"compression" = "ZSTD", -- 3. 고효율 압축
"colocate_with" = "sales_group" -- 5. Join 최적화
);
팁: 모든 컬럼에 인덱스를 걸면 데이터 적재(Load) 속도가 느려질 수 있으므로, 실제 쿼리에서 가장 많이 쓰이는 필터 조건 위주로 설정하는 것이 좋습니다.
References
• sql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md
• table_design/indexes/Bloomfilter_index.md
• table_design/data_compression.md
• best_practices/query_tuning/schema_tuning.md
References
• best_practices/partitioning.md
• best_practices/table_clustering.md
• table_design/table_design.md
• data_source/feature-support-data-lake-analytics.md
• sql-reference/data-types/semi_structured/VARIANT.md
• loading/loading_introduction/feature-support-loading-and-unloading.mdSangwon
06/04/2026, 6:17 AMRocky
06/04/2026, 6:17 AMstarrocks-spark-connector JAR 파일이 Spark 클래스패스에 포함되어 있어야 합니다.
* 네트워크: Spark 노드에서 StarRocks FE(8030, 9030) 및 BE(8040) 포트에 접근 가능해야 합니다.
2. PySpark 적재 예시
가장 범용적인 PySpark를 기준으로 한 코드 예시입니다.
python
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("StarRocks_Load").getOrCreate()
# 1. 적재할 데이터 프레임 생성 (예시)
df = spark.read.parquet("path/to/your/data")
# 2. StarRocks로 적재
df.write.format("starrocks") \
.option("starrocks.fe.http.url", "127.0.0.1:8030") \
.option("starrocks.fe.jdbc.url", "jdbc:<mysql://127.0.0.1:9030%22|mysql://127.0.0.1:9030">) \
.option("starrocks.table.identifier", "test_db.score_board") \
.option("starrocks.user", "root") \
.option("starrocks.password", "your_password") \
.mode("append") \
.save()
3. 주요 옵션 및 성능 팁
30만 건을 넣을 때 참고할 유용한 설정들입니다.
* mode("append") vs mode("overwrite"):
* `append`: 기존 데이터에 추가합니다. (Primary Key 테이블이면 키 중복 시 업데이트 발생)
* `overwrite`: 기존 데이터를 모두 지우고 새로 씁니다.
* 병렬 제어 (repartition):
* Spark의 파티션 개수만큼 StarRocks에 Stream Load 요청이 동시에 발생합니다.
* 30만 건 정도면 파티션이 너무 많을 필요가 없습니다. 만약 Spark 파티션이 수천 개라면 .repartition(5) 정도로 줄여서 StarRocks BE에 가해지는 동시 부하를 줄이는 것이 좋습니다.
* 트랜잭션 인터페이스 (v2.5 이상 권장):
* "starrocks.write.enable.transaction-stream-load" = "true" 설정을 주면 더 안정적이고 효율적으로 대량 적재를 수행할 수 있습니다. (최신 버전은 기본값이 TRUE인 경우가 많습니다.)
4. 주의사항
* Column Name: Spark DataFrame의 컬럼명과 StarRocks 테이블의 컬럼명이 일치해야 자동으로 매핑됩니다. (대소문자 구분 주의)
* JDBC 드라이버: Connector v1.1.1부터는 MySQL JDBC 드라이버가 포함되어 있지 않으므로, 별도로 `mysql-connector-java.jar`를 추가해야 합니다.
이 방식은 내부적으로 StarRocks의 Stream Load 메커니즘을 사용하기 때문에, 별도의 수동 파티셔닝이나 포맷 설정 없이도 매우 빠르게 데이터가 들어갑니다.
References
• loading/Spark-connector-starrocks.md
• unloading/Spark_connector.md
References
• loading/Spark-connector-starrocks.md
• unloading/Spark_connector.mdSangwon
06/04/2026, 7:03 AMRocky
06/04/2026, 7:03 AMpython
# 1. 원격지 테이블의 스키마 가져오기
table = get_table(spark, topic)
table_schema = table.schema
# 2. 시스템 컬럼(__로 시작하는 컬럼)을 제외하고 비즈니스 컬럼만 추출
target_fields = [f for f in table_schema.fields if not f.name.startswith("__")]
# 3. 필요한 컬럼만 선택하고 캐스팅(Cast) 수행
for field in target_fields:
# df에 해당 컬럼이 있는 경우에만 처리
if field.name in df.columns:
df = df.withColumn(field.name, df[field.name].cast(field.dataType))
# 4. 최종적으로 시스템 컬럼이 제거된, 타겟 테이블 구조에 맞는 컬럼들만 선택
final_column_names = [f.name for f in target_fields]
df = df.select(*final_column_names)
# 5. 이제 StarRocks로 Insert
# df.write.format("starrocks")...save()
왜 이런 컬럼이 보이나요?
1. 파티션 섀도잉: 원격지 테이블이 Iceberg, Hive 등 외부에 있고 특정 파티션 변환(Transformation)이 걸려있을 때, Spark가 이를 읽으면서 내부적으로 파티션 키를 계산하기 위해 임시 컬럼을 노출하는 경우가 있습니다.
2. StarRocks 내부 컬럼: StarRocks는 *op (Upsert/Delete 플래그)나 *row 같은 시스템 컬럼을 내부적으로 관리하는데, 특정 조회 방식(External Catalog 등)에서 이들이 스키마에 잡히기도 합니다.
추가 팁
* Case Sensitivity: StarRocks는 기본적으로 컬럼명의 대소문자를 구분하지 않는 설정이 많지만, Spark와 연동할 때는 가급적 *소문자로 통일*하거나 *테이블 정의와 정확히 일치*시키는 것이 안전합니다.
* Select 전용: 30만 건 정도면 `df.select()`를 통해 필요한 컬럼만 명시적으로 지정해서 넣는 것이 성능과 안정성 면에서 가장 깔끔합니다.
References
• sql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md (시스템 예약어 설명)
• loading/Spark-connector-starrocks.md (컬럼 매핑 가이드)
References
• sql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md
• data_source/catalog/iceberg/iceberg_meta_table.md
• sql-reference/information_schema/partitions.md
• sql-reference/information_schema/partitions_meta.md
• sql-reference/sql-statements/generated_columns.md