<@U0A71G31CDV> In StarRocks 4.0.4, I want to crea...
# questions-and-troubleshooting
s
@Rocky In StarRocks 4.0.4, I want to create a table within the internal catalog. The ⁠dt⁠ field is a string, and I'd like to configure it as a dynamic partition. I also want to apply the Duplicate Key model and 'DISTRIBUTE BY HASH'. Please provide an example.
r
In StarRocks 4.0.4, the most modern and recommended way to achieve "dynamic partitioning" for a
STRING
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
Copy code
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." Referencestable_design/data_distribution/expression_partitioning.mdadministration/management/FE_configuration.md
StarRocks 4.0.4에서
STRING
타입의 컬럼(
dt
)에 레거시 방식의 Dynamic Partition(동적 파티션) 설정을 적용하려면, 해당 문자열을 날짜로 변환하여 인식할 수 있도록
str2date
함수를 함께 사용해야 합니다. 다음은 요청하신 모든 설정값(
dynamic_partition
속성들)을 포함한 예시입니다. 예시: Duplicate Key 모델 및 동적 파티션 설정
Copy code
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
속성 방식이 이에 해당합니다.) Referencestable_design/data_distribution/dynamic_partitioning.mdsql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md Referencestable_design/data_distribution/dynamic_partitioning.mdtable_design/data_distribution/feature-support-data-distribution.md
s
@Rocky 파티션 테이블 조회방법은?
r
StarRocks에서 파티션 테이블을 조회하는 방법은 크게 세 가지 관점(데이터 조회, 파티션 메타 정보 조회, 쿼리 최적화 확인)으로 나눌 수 있습니다. 1. 일반적인 데이터 조회 (SELECT) 파티션 테이블이라고 해서 특별한 문법을 사용하는 것은 아닙니다. 일반 테이블과 동일하게
SELECT
문을 사용합니다.
Copy code
sql
SELECT * FROM my_table WHERE dt = '2024-01-01';
2. 파티션 프루닝(Partition Pruning) 활용 파티션 테이블을 조회할 때 가장 중요한 것은 *성능*입니다.
WHERE
절에 파티션 키(예:
dt
)를 조건으로 넣으면, StarRocks는 전체 데이터를 스캔하지 않고 해당 파티션의 데이터만 읽습니다. 이를 *파티션 프루닝*이라고 합니다. * 효율적인 조회: 파티션 컬럼을 조건에 포함하세요.
Copy code
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), 데이터 크기 등을 확인합니다.
Copy code
sql
    SHOW PARTITIONS FROM my_table;
* 시스템 뷰 조회 (v3.1.9+ / v3.2.4+): SQL을 통해 더 복잡한 조건으로 파티션을 조회할 수 있습니다.
Copy code
sql
    SELECT * FROM information_schema.partitions_meta
    WHERE TABLE_NAME = 'my_table'
    ORDER BY DATA_SIZE DESC;
4. 쿼리가 파티션을 잘 타는지 확인 (EXPLAIN) 내가 작성한 쿼리가 실제로 파티션 프루닝을 수행하는지 확인하려면 쿼리 앞에 `EXPLAIN`을 붙여 실행 계획을 확인하세요.
Copy code
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
수치 확인 Referencessql-reference/sql-statements/table_bucket_part_index/SHOW_PARTITIONS.mdbest_practices/partitioning.mdtable_design/data_distribution/expression_partitioning.md Referencessql-reference/sql-statements/table_bucket_part_index/SHOW_PARTITIONS.mdknowledge/admin_and_operation/shared_data_compaction.md
s
@Rocky 파티션을 여러개 하고싶은데, 아까전에 str2date 한 dt, layer_nm, fab_id 이렇게 세개가 가능해?
r
StarRocks 4.0.4 버전부터는 복합 표현식 파티셔닝(Mixed Expression Partitioning) 기능을 통해
dt
,
layer_nm
,
fab_id
세 개의 컬럼을 조합하여 파티션을 구성하는 것이 가능합니다. 가장 현대적이고 권장되는 방식인 *자동 파티셔닝(Automatic Partitioning)*을 사용하여 예시를 구성해 드립니다. 이 방식을 사용하면 데이터가 적재될 때
(날짜 + 레이어 + Fab)
조합에 맞는 파티션이 자동으로 생성됩니다. 여러 컬럼을 사용하는 파티션 테이블 예시
Copy code
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
절에 세 가지 조건을 모두 넣으면 성능이 극대화됩니다.
Copy code
sql
SELECT * FROM multi_partition_table 
WHERE dt = '2024-01-01' 
  AND layer_nm = 'PHOTO' 
  AND fab_id = 'FAB1';
이렇게 조회하면 StarRocks는 수많은 파티션 중 정확히 해당 조합의 파티션 하나만 읽으므로 매우 빠릅니다. Referencestable_design/data_distribution/expression_partitioning.mdtable_design/data_distribution/feature-support-data-distribution.mdsql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md Referencestable_design/data_distribution/feature-support-data-distribution.mdtable_design/data_distribution/Data_distribution.mdtable_design/data_distribution/list_partitioning.mdtable_design/data_distribution/dynamic_partitioning.md
s
@Rocky 파티션을 삭제하지않고 계속해서 유지하고싶어
r
데이터를 삭제하지 않고 영구히 보관하고 싶다면, 이전 답변에서 설명드린 설정 중 *데이터 유지 기간(TTL)과 관련된 속성들을 제거*하기만 하면 됩니다. StarRocks의 *자동 파티셔닝(Automatic Partitioning)*은 기본적으로 새로운 데이터가 들어올 때 파티션을 자동으로 생성하며, 별도의 삭제 설정을 하지 않으면 생성된 파티션은 사용자가 수동으로 삭제하기 전까지 영구히 유지됩니다. 영구 보관을 위한 파티션 테이블 설정 예시
Copy code
sql
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 테이블명`을 통해 파티션 개수를 확인해 보시는 것이 좋습니다. * 수동 삭제: 만약 나중에 특정 시점의 데이터가 정말 필요 없어져서 지우고 싶다면, 아래 명령어로 특정 파티션만 골라서 수동으로 삭제할 수 있습니다.
Copy code
sql
        ALTER TABLE permanent_multi_partition_table DROP PARTITION <파티션이름>;
결론적으로, `partition_live_number`나
dynamic_partition.start
같은 설정만 넣지 않으시면
데이터는 삭제되지 않고 계속 쌓이 유지됩니다. Referencestable_design/data_distribution/expression_partitioning.mdsql-reference/sql-statements/table_bucket_part_index/ALTER_TABLE.mdtable_design/data_distribution/dynamic_partitioning.md
s
@Rocky spark connector로 starrocks에 데이터 append하려고하는데, 테이블 스키마는 어떻게 가져와??
r
Spark Connector를 사용하여 StarRocks에 데이터를 Append할 때 스키마를 처리하는 방법은 크게 두 가지입니다. 1. 자동으로 스키마 가져오기 (권장) Spark Connector(v1.1.0 이상)는 내부적으로 *StarRocks의 FE JDBC 연결*을 통해 테이블의 메타데이터(컬럼명, 데이터 타입)를 자동으로 조회합니다. 따라서 별도로 스키마를 정의하지 않아도 Spark DataFrame의 컬럼명과 StarRocks의 컬럼명이 일치하면 자동으로 매핑됩니다. 이 기능을 사용하기 위해 설정에 반드시 *JDBC URL*이 포함되어야 합니다.
Copy code
scala
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 포맷 * 사용 예시:
Copy code
scala
    .option("starrocks.column.types", "user_id INT, dt DATE, revenue DECIMAL(10,2)")
3. 프로그램 코드 내에서 스키마 확인하기 현재 StarRocks 테이블의 스키마가 Spark에서 어떻게 인식되는지 미리 확인하고 싶다면, `load()`를 통해 빈 데이터를 읽어와서 출력해 볼 수 있습니다.
Copy code
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 파일을 반드시 클래스패스에 추가해야 합니다. Referencesloading/Spark-connector-starrocks.mdunloading/Spark_connector.md Referencesloading/Spark-connector-starrocks.mdunloading/Spark_connector.mdloading/Flink-connector-starrocks.mdunloading/Flink_connector.md Referencesloading/Spark-connector-starrocks.mdunloading/Spark_connector.md
s
@Rocky append하기전에 column명을 테이블 스키마와 일치하게 만들고싶은데 어떻게 하지?
r
Spark DataFrame의 컬럼명을 StarRocks 테이블 스키마와 일치시키는 가장 확실한 방법은 *StarRocks 테이블의 스키마 정보를 먼저 읽어온 뒤, 그 컬럼 리스트를 바탕으로 DataFrame을 변환*하는 것입니다. 다음은 PySpark와 Scala를 사용한 구체적인 방법입니다. 1. StarRocks에서 스키마 리스트 가져오기 StarRocks 커넥터의
load()
기능을 사용하면 실제 데이터를 읽지 않고도 스키마 정보만 메타데이터 형태로 가져올 수 있습니다. [PySpark 예시]
Copy code
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`를 사용해야 합니다.
Copy code
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와 달라도 매핑이 가능합니다.
Copy code
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)`를 적용하는 것이 가장 자동화된 방법입니다. Referencesloading/Spark-connector-starrocks.mdunloading/Spark_connector.mdsql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md
s
@Rocky 캐스팅은 어떻게해??
r
Spark Connector를 사용하여 StarRocks에 데이터를 넣을 때 *캐스팅(데이터 타입 변환)*은 크게 세 가지 방법으로 할 수 있습니다. 1. Spark DataFrame에서 직접 캐스팅 (가장 권장) 데이터를
write
하기 전에 Spark의
cast()
함수를 사용하여 StarRocks 테이블의 데이터 타입과 일치시키는 방법입니다. Spark 수준에서 데이터를 검증할 수 있어 가장 안전합니다.
Copy code
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
)를 거쳐야 하는 경우 유용합니다. * 예시:
Copy code
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 같은 특수 타입을 처리할 때 주로 사용합니다.
Copy code
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()`를 통해 타입을 명확히 맞춘 후 전송하는 것이 디버깅에 유리합니다. Referencesloading/Spark-connector-starrocks.mdunloading/Spark_connector.mdsql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md Referencesunloading/Spark_connector.mdloading/Spark-connector-starrocks.mdunloading/Flink_connector.md
s
@Rocky spark connector로 데이터 넣을때, upsert는 어떻게해
r
StarRocks Spark Connector를 통해 *Upsert(데이터가 있으면 Update, 없으면 Insert)*를 수행하려면, 대상 테이블이 반드시 *Primary Key(PK) 모델*이어야 합니다. 1. 기본 Upsert (전체 컬럼 업데이트) 별도의 설정 없이 PK 테이블에 데이터를
append
모드로 쓰면 자동으로 Upsert가 수행됩니다. PK가 동일한 행이 들어오면 기존 데이터를 덮어쓰고, 없으면 새로 추가합니다.
Copy code
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 포함 필수)
Copy code
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`: 기준이 될 컬럼명 지정
Copy code
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`만으로 충분합니다. Referencesloading/Spark-connector-starrocks.mdloading/Load_to_Primary_Key_tables.mdtable_design/table_types/primary_key_table.md Referencesloading/Spark-connector-starrocks.mdloading/Load_to_Primary_Key_tables.mdtable_design/table_types/primary_key_table.mdloading/Flink-connector-starrocks.mdunloading/Spark_connector.md Referencesloading/Spark-connector-starrocks.mdloading/Load_to_Primary_Key_tables.mdtable_design/table_types/primary_key_table.mdloading/Flink-connector-starrocks.mdunloading/Spark_connector.mdunloading/Flink_connector.md Referencesloading/Spark-connector-starrocks.mdloading/Kafka-connector-starrocks.md
s
@Rocky primary_key, partiton by dt, layer,fab 이렇게해서 예제 하나만 줘
r
StarRocks에서 *Primary Key 모델*을 사용하면서
dt
,
layer
,
fab
세 가지 컬럼으로 파티셔닝을 하는 예제입니다. 주의사항 1. PK 제약 조건: Primary Key 모델에서 파티션 키(`PARTITION BY`에 사용되는 컬럼)는 반드시 *Primary Key 컬럼 구성 요소에 포함*되어야 합니다. 2. 자동 파티션 생성:
dt
,
layer
,
fab
조합이 많을 경우 일일이 파티션을 지정하기 힘들기 때문에, List 자동 파티션(Automatic Partitioning) 기능을 사용하는 것이 효율적입니다. 테이블 생성 SQL 예제
Copy code
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 파티션) 특정 값들을 미리 지정해야 한다면 아래와 같은 형식을 사용합니다.
Copy code
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
조합은 가변적이므로 첫 번째 예제인 자동 파티션 방식을 가장 추천합니다.* Referencessql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.mdtable_design/table_types/primary_key_table.mdtable_design/partitioning/automatic_partitioning.md
s
@Rocky The partiton column could not be aggregated column and unique table’s partition column must be key column 이런?
r
네, 정확히 보셨습니다! 말씀하신 내용은 StarRocks의 테이블 모델별 제약 사항입니다. StarRocks에서는 데이터 무결성과 파티션 프루닝(Partition Pruning)을 위해 *파티션 키(Partition Key) 설정에 제약*을 두고 있습니다. 각 모델별 규칙은 다음과 같습니다. 1. Primary Key 테이블 (가장 엄격) * 규칙: 파티션 컬럼은 반드시 *Primary Key의 일부*여야 합니다. * 이유: Primary Key 모델은 쓰기 시점에 중복 체크를 수행하는데, 파티션 키가 PK에 포함되어야 어느 파티션에서 중복을 확인해야 할지 결정할 수 있기 때문입니다. 2. Unique Key 테이블 * 규칙: 파티션 컬럼은 반드시 *Key 컬럼(Unique Key)*에 포함되어야 합니다. * 이유: Unique 모델도 동일 키 데이터가 들어오면 업데이트(Replace)를 해야 하므로, 해당 키 데이터가 저장될 파티션을 명확히 찾아야 합니다. 3. Aggregate 테이블 (집계 테이블) * 규칙: 파티션 컬럼은 반드시 *Key 컬럼*이어야 하며, Value 컬럼(집계 함수가 적용된 컬럼)은 파티션 컬럼이 될 수 없습니다. * 이유: Aggregate 모델은 Key가 같은 데이터를 합치는 구조입니다. 집계가 수행될 결과값(SUM, REPLACE 등)을 기준으로 데이터를 나누는 것은 논리적으로 불가능합니다. 4. Duplicate Key 테이블 (가장 유연) * 규칙: 어떤 컬럼이든 파티션 컬럼으로 사용할 수 있습니다. (Key가 아니어도 됨) * 이유: 단순히 데이터를 쌓는 구조이므로 중복 체크나 집계 제약이 없기 때문입니다. ────────── 요약 예시 (Primary Key 테이블 기준) 질문하셨던
dt
,
layer
,
fab
예제를 다시 보면:
Copy code
sql
PRIMARY KEY(dt, layer, fab, id) -- 이 4개 조합이 PK
PARTITION BY (dt, layer, fab)   -- PK 구성 요소 중 일부이므로 '가능'
만약 아래처럼 선언하면 *에러*가 발생합니다.
Copy code
sql
PRIMARY KEY(id)                 -- PK는 id 하나뿐인데
PARTITION BY (dt)               -- PK에 없는 dt로 파티션을 나누려 하면 '에러'
따라서 PK 테이블에서 파티션을 나누고 싶다면, *해당 파티션 컬럼들을 모두 PK 선언부(괄호 안)에 포함*시켜 주어야 합니다. Referencestable_design/partitioning/table_partitioning.mdtable_design/table_types/primary_key_table.mdsql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md
s
@Rocky duplicate key로 하니까, append할때, 데이터 중복이 되던데 이건 upsert로 동작하지않아?
r
아니요, Duplicate Key 모델은 Upsert를 지원하지 않습니다. Duplicate Key 모델은 이름 그대로 "중복을 허용하는" 모델입니다. Spark에서
append
모드로 데이터를 넣으면, 기존에 동일한 ID나 키가 있더라도 상관하지 않고 *매번 새로운 행으로 추가*합니다. 원하시는 Upsert(있으면 Update, 없으면 Insert) 기능을 사용하시려면 반드시 *Primary Key 모델*을 사용해야 합니다. 왜 Duplicate Key에서는 안 되나요? * Duplicate Key 모델: 데이터가 들어오는 대로 단순히 쌓는 구조입니다. 분석용 원천 로그 데이터를 그대로 저장할 때 사용하며, '업데이트'라는 개념 자체가 없습니다. * Primary Key 모델: 데이터를 넣을 때 메모리에 있는 인덱스(Persistent Index)를 확인하여, 동일한 키가 있으면 해당 위치의 데이터를 수정하고 없으면 새로 넣습니다. 해결 방법: 테이블 모델 변경 테이블을 생성할 때 아래와 같이
PRIMARY KEY
모델로 만드셔야 합니다.
Copy code
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가 불가능한 구조이므로 테이블 정의를 변경하셔야 합니다. Referencestable_design/table_types/duplicate_key_table.mdtable_design/table_types/primary_key_table.mdloading/Spark-connector-starrocks.md
s
@Rocky duplicated key 테이블이 뭔지 간략하게설명햐줘
r
*Duplicate Key 테이블*은 StarRocks의 가장 기본적이고 단순한 테이블 모델로, 한마디로 *"데이터를 중복에 상관없이 들어오는 대로 쌓는 모델"*입니다. 핵심 특징은 다음과 같습니다. 1. 중복 허용 (No Deduplication) * 동일한 데이터를 여러 번 입력해도 기존 데이터를 덮어쓰거나 합치지 않고, *모든 행을 그대로 저장*합니다. * 따라서 `Upsert`(업데이트)가 발생하지 않으며, 모든 데이터는 `Append`(추가)로 처리됩니다. 2. 정렬 키 (Sorting Key)의 역할 * `DUPLICATE KEY(dt, id)`와 같이 키를 지정하는 이유는 *"이 컬럼들을 기준으로 데이터를 정렬해서 저장해라"*라는 뜻입니다. * 데이터가 정렬되어 저장되므로, 해당 컬럼을 조건(
WHERE dt = '...'
)으로 조회할 때 검색 속도가 매우 빠릅니다. (유니크 제약과는 무관합니다.) 3. 가장 빠른 로딩 속도 * 새 데이터를 넣을 때 기존 데이터와 비교하거나 중복 체크를 하는 과정이 전혀 없습니다. * 따라서 4가지 테이블 모델(Duplicate, Primary, Unique, Aggregate) 중 쓰기(Ingestion) 성능이 가장 뛰어납니다. 4. 언제 사용하나요? * 원천 로그 데이터: 클릭 로그, 센서 데이터처럼 중복 삭제가 필요 없고 발생한 그대로를 저장해야 하는 경우. * 단순 적재: 데이터 수정이 거의 없고, 오로지 빠르게 쌓고 조회(Filter)하는 것이 중요한 경우. 요약 예시
Copy code
sql
CREATE TABLE logs (
    dt DATE,
    id INT,
    event VARCHAR(20)
) 
DUPLICATE KEY(dt, id) -- 이 컬럼들로 정렬만 수행함
...
위 테이블에 `(2023-10-01, 1, 'click')`을 두 번 INSERT하면, 테이블에는 똑같은 데이터가 2줄 남게 됩니다. Referencestable_design/table_types/duplicate_key_table.mdtable_design/table_types/table_types.md
s
@Rocky starrocks internal 로 table create 할때, 성능을 위해 properties에 parquet파일로 만든다거나, 넣을떄 정렬조건을 준다던가 bloom-filter를 준다는 등의 유용한 설정들도 줄수있을까? 있으면 추천해줘
r
StarRocks 내부(Internal) 테이블을 생성할 때 성능을 극대화할 수 있는 주요 설정들을 추천해 드립니다. StarRocks 내부 테이블은 기본적으로 자체 최적화된 Columnar Storage(Segment V2) 포맷을 사용하므로, 사용자가 직접 "Parquet"으로 포맷을 지정할 필요는 없습니다. 대신 아래의 설정들을 통해 조회 및 저장 성능을 튜닝할 수 있습니다. ────────── 1. 정렬 키 (Sort Key) 설정 StarRocks는 데이터를 정렬해서 저장하며, 이를 통해 *Prefix Index(앞부분 일치 인덱스)*와 Zone-map(Min/Max 필터링) 효과를 얻습니다. * 설정 방법:
DUPLICATE KEY
,
PRIMARY KEY
,
UNIQUE KEY
등에 정의된 컬럼 순서대로 데이터가 정렬됩니다. * 추천:
WHERE
절에 가장 자주 쓰이는 컬럼, 변별력(Cardinality)이 높은 컬럼을 *가장 앞쪽*에 배치하세요. 2. Bloom Filter 인덱스 (
bloom_filter_columns
)
특정 값이 데이터에 포함되어 있는지 빠르게 확인하여 불필요한 파일을 읽지 않게 합니다. * 용도: `id`나
serial_number
같이 값의 종류가 매우 많고(High Cardinality),
=
또는
IN
조건으로 자주 조회하는 컬럼에 효과적입니다. * 제약:
TINYINT
,
FLOAT
,
DOUBLE
,
DECIMAL
타입은 지원하지 않습니다. * 예시:
Copy code
sql
    PROPERTIES (
        "bloom_filter_columns" = "user_id, order_id"
    )
3. 데이터 압축 방식 (
compression
)
저장 공간을 아끼고 I/O를 줄여 성능을 높입니다. * 추천: *
LZ4
(기본값): 압축/해제 속도가 매우 빨라 CPU 부하가 적습니다. * `ZSTD`: 압축률이 매우 좋아 저장 공간을 크게 절약하며, I/O가 병목인 환경에서 성능이 더 좋을 수 있습니다. * 예시:
Copy code
sql
    PROPERTIES (
        "compression" = "ZSTD"
    )
4. Inverted Index (역색인, v3.3+ 권장) 텍스트 검색이나 복합적인 필터링 성능을 획기적으로 높여줍니다. * 용도: 로그 데이터의 텍스트 검색이나 여러 컬럼을 조합한 필터링이 많을 때 유리합니다. * 예시:
Copy code
sql
    INDEX index_name (column_name) USING INVERTED [COMMENT 'comment']
5. Colocate Join (
colocate_with
)
자주 Join하는 두 테이블의 데이터를 물리적으로 동일한 노드에 배치하여, 조인 시 네트워크 전송(Shuffle)을 없앱니다. * 설정 방법: Join 대상 테이블들에 동일한 그룹명을 지정합니다. * 예시:
Copy code
sql
    PROPERTIES (
        "colocate_with" = "group_name_1"
    )
6. Bitmap Index 값의 종류가 적은(Low Cardinality) 컬럼에 대해 필터링 속도를 높입니다. * 용도:
gender
,
status
,
region
등 값의 종류가 몇 개 안 되는 컬럼에 적합합니다. ────────── 종합 추천 예시 SQL
Copy code
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) 속도가 느려질 수 있으므로, 실제 쿼리에서 가장 많이 쓰이는 필터 조건 위주로 설정하는 것이 좋습니다. Referencessql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.mdtable_design/indexes/Bloomfilter_index.mdtable_design/data_compression.mdbest_practices/query_tuning/schema_tuning.md Referencesbest_practices/partitioning.mdbest_practices/table_clustering.mdtable_design/table_design.mddata_source/feature-support-data-lake-analytics.mdsql-reference/data-types/semi_structured/VARIANT.mdloading/loading_introduction/feature-support-loading-and-unloading.md
s
@Rocky Spark connector로 starrocks internal catalog에 30만건을 넣어야하는데, 어떻게 넣어야해?
r
Spark Connector를 사용하여 StarRocks 내부(Internal) 테이블에 30만 건의 데이터를 적재하는 방법은 간단합니다. 30만 건 정도의 데이터는 StarRocks 입장에서 매우 적은 양이므로, 기본적인 Batch Write 방식으로 충분히 빠르게 처리할 수 있습니다. 1. 필수 설정 (Prerequisites) * JAR 파일:
starrocks-spark-connector
JAR 파일이 Spark 클래스패스에 포함되어 있어야 합니다. * 네트워크: Spark 노드에서 StarRocks FE(8030, 9030) 및 BE(8040) 포트에 접근 가능해야 합니다. 2. PySpark 적재 예시 가장 범용적인 PySpark를 기준으로 한 코드 예시입니다.
Copy code
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 메커니즘을 사용하기 때문에, 별도의 수동 파티셔닝이나 포맷 설정 없이도 매우 빠르게 데이터가 들어갑니다. Referencesloading/Spark-connector-starrocks.mdunloading/Spark_connector.md Referencesloading/Spark-connector-starrocks.mdunloading/Spark_connector.md
s
@Rocky Insert하기전에 원격지에 있는 컬럼을 조회하면 __generated_partition_column_0이라고 나오거든?? 이걸빼면 되지않을까?? 지금 코드를 어떻게 짰냐면 Table = get_table(spark, topic) Table_schema = table.schema For field in table_schema.fields: Df = df.withColumn(field.name, df[field.name].cast(field.dataType))
r
네, 맞습니다. *`*generated_partition_column_0`과 같이 `*`로 시작하는 컬럼들은 StarRocks나 Spark가 내부적으로 사용하는 시스템/메타데이터 컬럼*입니다. 이 컬럼들은 실제 데이터 테이블의 비즈니스 컬럼이 아니라, 파티션 관리나 내부 연산을 위해 생성된 "가상 컬럼"일 가능성이 매우 높습니다. 따라서 *내부(Internal) 테이블로 Insert할 때는 이 컬럼들을 제외하는 것이 정답*입니다. 제외하지 않으면 StarRocks가 "해당 컬럼이 테이블 정의에 없다"는 에러를 내거나, 불필요한 데이터를 저장하려고 시도하게 됩니다. 제공해주신 코드를 다음과 같이 수정하여 *시스템 컬럼을 필터링*하는 로직을 추가하는 것을 추천합니다. 수정된 코드 제안
Copy code
python
# 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()`를 통해 필요한 컬럼만 명시적으로 지정해서 넣는 것이 성능과 안정성 면에서 가장 깔끔합니다. Referencessql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.md (시스템 예약어 설명) • loading/Spark-connector-starrocks.md (컬럼 매핑 가이드) Referencessql-reference/sql-statements/table_bucket_part_index/CREATE_TABLE.mddata_source/catalog/iceberg/iceberg_meta_table.mdsql-reference/information_schema/partitions.mdsql-reference/information_schema/partitions_meta.mdsql-reference/sql-statements/generated_columns.md