<@U0A71G31CDV> Iceberg+Minio에 넣으려고하는데, Starrocks S...
# questions-and-troubleshooting
s
@Rocky Iceberg+Minio에 넣으려고하는데, Starrocks Stream_load를 사용해서 external catalog에 넣을수있어?
r
StarRocks의 *Stream Load*는 기본적으로 StarRocks의 *Internal Table(내부 테이블)*에 데이터를 적재하기 위한 도구입니다. 현재 버전 기준으로 External Catalog(Iceberg 등)의 외부 테이블에 직접 `STREAM LOAD`를 사용하는 것은 지원되지 않습니다. 대신, MinIO의 Iceberg 테이블에 데이터를 넣으려면 다음과 같은 방법을 사용해야 합니다. 1.
INSERT INTO
구문 사용 (추천)
StarRocks v3.1부터 Iceberg 테이블에 대한 `INSERT INTO`를 지원합니다. * 파일이 이미 MinIO에 있는 경우:
FILES()
함수를 사용하여 직접 Iceberg 테이블로 넣을 수 있습니다.
Copy code
sql
    INSERT INTO iceberg_catalog.db_name.table_name
    SELECT * FROM FILES(
        "path" = "<s3://bucket/data.parquet>",
        "format" = "parquet",
        "aws.s3.endpoint" = "<http://minio:9000>",
        "aws.s3.access_key" = "admin",
        "aws.s3.secret_key" = "password"
    );
* 직접 값을 넣는 경우:
Copy code
sql
    INSERT INTO iceberg_catalog.db_name.table_name VALUES (1, 'data');
2. 로컬 파일을 꼭 Stream Load로 넣어야 하는 경우 (2단계 방식) 로컬에 있는 데이터를 반드시 전송해야 한다면, 중간 임시 테이블을 거쳐야 합니다. 1. StarRocks 내부에 임시 테이블(OLAP 테이블)을 생성합니다. 2. `STREAM LOAD`를 사용하여 데이터를 임시 테이블에 넣습니다. 3.
INSERT INTO iceberg_table SELECT * FROM temp_table
명령으로 Iceberg 테이블에 적재합니다. 요약 및 주의사항 * 지원 버전: Iceberg 테이블에 데이터를 쓰는(Sink) 기능은 *v3.1 이상*부터 지원됩니다. * 포맷: 현재 Iceberg 외부 테이블로 데이터를 쓸 때는 Parquet 포맷만 지원됩니다. * 권한: Catalog 생성 시 S3(MinIO)에 대한 쓰기 권한이 설정되어 있어야 합니다. 따라서 Stream Load API를 직접 Iceberg 테이블 엔드포인트로 쏘는 것은 불가능하며, *`INSERT INTO`*가 공식적인 데이터 연동 방법입니다. Referencesloading/StreamLoad.mdsql-reference/sql-statements/loading_unloading/STREAM_LOAD.md
s
@Rocky 그럼 다음으로 Internal Catalog에 데이터를 넣을때 가장 좋은 방법 두 가지 정도를 추천해줘 한번에 30만건의 데이터를 넣을거야
r
StarRocks의 *Internal Catalog(내부 테이블)*에 30만 건 정도의 데이터를 한 번에 넣을 때 가장 효율적이고 권장되는 방법 2가지는 다음과 같습니다. StarRocks는 OLAP 엔진이므로 "한 번에(Batch)" 넣는 것이 성능과 시스템 안정성 면에서 가장 중요합니다. 30만 건 정도는 StarRocks에게 매우 가벼운 수준이며, 아래 방법으로 수 초 내에 처리가 가능합니다. ────────── 1. Stream Load (가장 빠름, 로컬/애플리케이션 데이터용) 데이터가 현재 로컬 파일(CSV, JSON, Parquet 등)로 있거나, 프로그램(Python, Java 등)에서 생성된 데이터라면 *Stream Load*가 가장 좋습니다. * 특징: HTTP 프로토콜을 사용하며, StarRocks BE에 직접 데이터를 쏴주는 방식입니다. 오버헤드가 가장 적고 빠릅니다. * 실행 예시 (curl 기준):
Copy code
bash
    curl --location-trusted -u <username>:<password> \
        -H "label:load_300k_rows" \
        -H "format: csv" \
        -H "column_separator: ," \
        -T data.csv \
        http://<fe_host>:<fe_http_port>/api/<db_name>/<table_name>/_stream_load
* 장점: 별도의 스토리지 거칠 필요 없이 즉시 적재 가능. 2. INSERT INTO SELECT FROM FILES() (MinIO/S3 데이터용) 데이터가 이미 MinIO(또는 S3/HDFS)와 같은 오브젝트 스토리지에 올라와 있다면, *`FILES()` 테이블 함수*를 사용하는 것이 가장 편리합니다. * 특징: SQL 구문만으로 외부 저장소의 파일을 직접 읽어 내부 테이블로 복사합니다. StarRocks v3.1 이상에서 권장하는 현대적인 방식입니다. * 실행 예시 (SQL):
Copy code
sql
    INSERT INTO internal_db.target_table
    SELECT * FROM FILES(
        "path" = "<s3://my-bucket/data/300k_records.parquet>",
        "format" = "parquet",
        "aws.s3.endpoint" = "<http://minio:9000>",
        "aws.s3.access_key" = "admin",
        "aws.s3.secret_key" = "password"
    );
* 장점: 별도의 로드 도구가 필요 없고 SQL만으로 제어 가능하며, 병렬 처리 성능이 매우 뛰어남. ────────── 💡 주의 및 팁 * 절대 피해야 할 방법: 30만 건의 데이터를
INSERT INTO VALUES (1, ...), (2, ...)
식의 수많은 개별 쿼리로 나누어 던지는 것은 지양해야 합니다. (StarRocks의 Transaction 오버헤드가 발생하고 데이터 버전이 너무 많아집니다.) * Batch 처리: 위 2가지 방법 모두 30만 건을 *하나의 트랜잭션*으로 묶어서 처리하므로 StarRocks의 성능을 최대로 활용할 수 있습니다. * 권장사항: 데이터가 로컬에 있다면 1번(Stream Load), MinIO에 있다면 *2번(FILES)*을 선택하세요. Referencesloading/StreamLoad.mdloading/InsertInto.mdsql-reference/sql-functions/table-functions/files.md Referencesknowledge/usage/best-practice.mdloading/InsertInto.mdloading/Loading_intro.mdloading/loading_introduction/loading_considerations.md
s
@Rocky 그럼 pandas dataframe을 streamload류 넣를 수잏어?? Spark connector가 공식문서에서 추천이던데 그거보다 나은점이 뭐길래추천하는거야
r
네, Pandas DataFrame을 Stream Load를 통해 StarRocks에 넣을 수 있습니다. 또한 Spark Connector가 왜 추천되는지, Stream Load와 비교했을 때 어떤 장점이 있는지 설명해 드릴게요. 1. Pandas DataFrame을 Stream Load로 넣는 방법 Pandas에는 전용 "StarRocks 커넥터"가 따로 있는 것은 아니지만, *Stream Load가 HTTP 기반 API*이기 때문에 Python의
requests
라이브러리를 사용하여 간단히 구현할 수 있습니다. * 방법: DataFrame을 메모리 상에서 CSV 또는 Parquet 버퍼로 변환한 뒤, 이를 Stream Load 엔드포인트로 전송합니다. * 예시 코드 (Python):
Copy code
python
    import pandas as pd
    import requests
    from io import BytesIO

    df = pd.DataFrame(...) # 30만 건의 데이터

    # 1. DataFrame을 CSV 버퍼로 변환
    csv_buffer = BytesIO()
    df.to_csv(csv_buffer, index=False, header=False)
    csv_buffer.seek(0)

    # 2. Stream Load 호출
    url = "http://<FE_HOST>:<FE_HTTP_PORT>/api/<DB>/<TABLE>/_stream_load"
    auth = ("root", "your_password")
    headers = {
        "label": "label_12345",
        "format": "csv",
        "column_separator": ","
    }

    response = requests.put(url, data=csv_buffer, auth=auth, headers=headers)
    print(response.json())
30만 건 정도는 이 방식으로도 수 초 내에 매우 빠르게 처리됩니다. ────────── 2. Spark Connector를 추천하는 이유 (vs Stream Load 직접 사용) 공식 문서에서 Spark Connector를 추천하는 이유는 *"대규모 분산 처리"*와 "편의성" 때문입니다. 사실 *Spark Connector도 내부적으로는 Stream Load 기술을 사용*합니다. 직접 Stream Load를 쓰는 것보다 나은 점은 다음과 같습니다: 1. 분산 처리 (Scalability): * Pandas는 단일 노드(메모리)에서 동작하지만, Spark는 여러 워커 노드에서 데이터를 병렬로 읽고 StarRocks의 여러 BE(Backend) 노드에 동시에 적재합니다. TB 단위의 대용량 데이터라면 Spark가 필수적입니다. 2. 내결함성 (Fault Tolerance): * 적재 중 실패했을 때 체크포인트를 통한 재시도, 트랜잭션 관리 등을 커넥터가 자동으로 처리해 줍니다. 3. 데이터 타입 매핑 자동화: * 직접 Stream Load를 쓸 때는 CSV 헤더나 컬럼 순서를 수동으로 맞춰야 하지만, Spark Connector는 DataFrame의 스키마를 보고 StarRocks 테이블 구조에 맞춰 자동으로 매핑해 줍니다. 4. 복잡한 ETL 연동: * 데이터를 단순히 넣는 게 아니라, Spark에서 조인(Join)이나 집계(Aggregation)를 마친 결과를 바로 StarRocks에 넣고 싶을 때 파이프라인 구성이 훨씬 깔끔합니다. 결론: 무엇을 써야 할까요? * 현재 데이터가 Pandas DataFrame에 있고, 크기가 30만~100만 건 정도라면: 굳이 무거운 Spark를 띄울 필요 없이 *Stream Load(Python requests 방식)*를 추천합니다. 훨씬 가볍고 빠릅니다. * 데이터가 이미 HDFS/S3에 있거나, 수천만 건 이상의 대용량이며, 기존에 Spark 환경을 사용 중이라면: *Spark Connector*를 사용하는 것이 훨씬 안정적이고 관리하기 편합니다. Referencesloading/Spark-connector-starrocks.mdloading/StreamLoad.mdloading/Loading_intro.md Referencesloading/Spark-connector-starrocks.mdloading/Loading_intro.md
s
@Rocky no vliad Basic Authorization이라는데?
r
no valid Basic Authorization
에러가 발생하는 이유는 StarRocks의 *Stream Load 작동 방식(Redirection)*과 Python
requests
라이브러리의 보안 정책 때문입니다. 1. 에러가 발생하는 원인 1. Redirection: 사용자가 FE(Front-End) 노드로 Stream Load를 요청하면, FE는 실제 데이터를 처리할 BE(Back-End) 노드로 요청을 보냅니다 (HTTP 307 Redirect). 2. Auth Header 제거: Python의
requests
라이브러리는 보안상의 이유로 *서로 다른 호스트(FE -> BE)로 리다이렉트될 때
Authorization
헤더를 삭제*하고 전송합니다. 그래서 BE 노드에서 "인증 정보가 없다"는 에러가 나는 것입니다. 2. 해결 방법 (Python 코드 수정) `requests`가 헤더를 지우지 못하게 하거나, 리다이렉트를 수동으로 처리해야 합니다. 가장 간단한 방법은 *헤더를 직접 만들어서 넣고, 리다이렉트를 수동으로 처리*하는 것입니다.
Copy code
python
import requests
import base64

def starrocks_stream_load(df, host, port, db, table, user, password):
    # 1. Basic Auth 헤더 직접 생성
    auth_str = f"{user}:{password}"
    encoded_auth = base64.b64encode(auth_str.encode()).decode()
    
    url = f"http://{host}:{port}/api/{db}/{table}/_stream_load"
    headers = {
        "Authorization": f"Basic {encoded_auth}",
        "Expect": "100-continue",
        "label": "my_label_12345",
        "format": "csv",
        "column_separator": ","
    }

    # 2. 첫 요청: FE에 전송 (리다이렉트 자동 방지)
    # 30만 건 정도면 CSV 스트링으로 변환해서 보냅니다.
    csv_data = df.to_csv(index=False, header=False)
    
    # allow_redirects=True로 하면 헤더가 증발하므로 주의
    # 가장 확실한 방법은 리다이렉트 응답을 받아서 직접 BE로 쏘는 것입니다.
    res = requests.put(url, data=csv_data, headers=headers, allow_redirects=True)
    
    # 만약 위 코드로도 안 된다면, 아래처럼 Session을 커스텀하거나 
    # 혹은 직접 BE 노드 IP/Port(8040)로 바로 쏘는 방법이 있습니다.
    return res.json()
────────── 3. Spark Connector를 추천하는 이유 (vs Stream Load) 공식 문서에서 Spark Connector를 강력히 추천하는 이유는 "대용량 병렬 처리 최적화" 때문입니다. * 분산 적재 (Parallel Loading): Stream Load는 보통 하나의 클라이언트가 하나의 BE에 데이터를 쏘지만, Spark Connector는 Spark의 워커 노드 수만큼 여러 BE에 동시에 데이터를 나눠서 쏩니다. * Transaction 관리: Spark 커넥터는 내부적으로 2-Phase Commit 같은 메커니즘을 지원하여, 도중에 일부가 실패했을 때 데이터 정합성을 맞추는 처리가 더 견고합니다. * 데이터 셔플링 최적화: StarRocks의 데이터 분포(Bucketing)에 맞춰 Spark 단계에서 미리 데이터를 분류(Shuffle)한 뒤 적재하므로, StarRocks BE 내부의 부하를 획기적으로 줄여줍니다. * 복잡한 타입 처리: Pandas 등에서 발생할 수 있는 데이터 타입 불일치 문제를 Spark 스키마 기반으로 자동 변환해 줍니다. 결론: * 30만 건 정도의 가벼운 작업: Python + Stream Load로 충분합니다. 위에서 말한 인증 헤더 문제만 해결하면 됩니다. * 수천만 건 ~ TB 단위 또는 정기적인 ETL: 무조건 *Spark Connector*를 쓰는 것이 성능과 운영 안정성 면에서 훨씬 유리합니다. Referencessql-reference/sql-statements/loading_unloading/STREAM_LOAD.mdloading/Spark-connector-starrocks.mdfaq/loading/Stream_load_faq.md Referencessql-reference/sql-statements/loading_unloading/STREAM_LOAD.mdloading/Stream_Load_transaction_interface.md