LLM WikiAccess-protected knowledge portal
← 스터디 홈
2편 · 약 22분

Trino 커넥터 심화: Iceberg, MySQL, Kafka 커넥터 운영

커넥터는 Trino의 데이터 접속 계층이다

1편에서 커넥터 프레임워크의 구조를 살펴봤다. 이번 편에서는 실제 데이터 플랫폼에서 가장 자주 마주치는 세 커넥터를 깊이 다룬다.

  • Iceberg 커넥터: S3나 HDFS 위에 올린 Lakehouse 레이어. 스냅샷, 파티션 프루닝, 스키마 발전(evolution).
  • MySQL 커넥터: 운영 DB에서 직접 데이터를 읽거나 Trino를 통한 크로스-소스 조인.
  • Kafka 커넥터: 실시간 토픽을 SQL로 조회. 스트림-배치 브리징.

각 커넥터의 설정, 한계, 운영 주의점을 구체적으로 짚는다.


Iceberg 커넥터

카탈로그 타입 선택

Iceberg는 테이블 메타데이터를 별도의 카탈로그에 저장한다. Trino의 iceberg.catalog.type 속성으로 어떤 카탈로그를 쓸지 지정한다.

Iceberg 카탈로그 옵션 비교 HMS (Hive Metastore) • 기존 Hadoop 환경 • Hive와 테이블 공유 가능 • Thrift RPC 프로토콜 • 운영 오버헤드 있음 AWS Glue • AWS 환경 권장 • 서버리스, 관리형 • IAM 통합 • API 호출 비용 발생 REST Catalog • 클라우드 네이티브 • Nessie·Polaris·Unity • OAuth2 인증 지원 • View 관리 미지원 JDBC / Nessie • 소규모/테스트 환경 • Nessie = Git-like 버전 관리 가능 • 프로덕션엔 제한적 오브젝트 스토리지 (S3 / GCS / ADLS) 실제 Parquet / ORC 파일 + Iceberg 메타데이터 JSON 카탈로그는 메타데이터 포인터만 관리 — 실제 파일은 항상 오브젝트 스토리지에
Iceberg 카탈로그 타입 선택 기준

HMS 카탈로그 설정

# etc/catalog/iceberg.properties
connector.name=iceberg
iceberg.catalog.type=hive_metastore
hive.metastore.uri=thrift://hms-host:9083
hive.metastore.authentication.type=NONE

hive.metastore.authentication.type=KERBEROS로 Kerberized 환경도 지원한다.

REST 카탈로그 설정

connector.name=iceberg
iceberg.catalog.type=rest
iceberg.rest-catalog.uri=http://iceberg-rest-catalog:8181
# OAuth2 인증이 필요한 경우
iceberg.rest-catalog.security=OAUTH2
iceberg.rest-catalog.oauth2.token=<token>

REST 카탈로그는 Polaris(Apache), Unity Catalog(Databricks Open Source), Nessie 같은 현대적인 메타데이터 서버와 함께 사용한다. 단, 2026년 기준 Trino의 REST 카탈로그는 View 생성/관리를 지원하지 않는다는 점을 알아야 한다.

Iceberg 커넥터 핵심 기능

#### 타임 트래블(Time Travel)

Iceberg 스냅샷을 활용해 과거 데이터를 조회할 수 있다.

-- 특정 스냅샷 ID로 조회
SELECT * FROM orders FOR VERSION AS OF 5432123456789;

-- 특정 시점 기준 조회
SELECT * FROM orders FOR TIMESTAMP AS OF TIMESTAMP '2026-01-15 00:00:00';

이를 활용하면 파이프라인 오류 후 이전 상태의 데이터로 재처리할 수 있다.

#### 파티션 프루닝(Partition Pruning)

Iceberg는 데이터 파일마다 파티션 통계(min/max, null count)를 메타데이터에 저장한다. WHERE 조건이 파티션 컬럼과 일치하면 Trino가 불필요한 파일 읽기를 건너뛴다.

-- partition_date가 파티션 컬럼이면, 해당 날짜의 파일만 읽음
SELECT * FROM events WHERE partition_date = DATE '2026-07-01';

파티션 프루닝이 동작하려면 WHERE 조건이 파티션 컬럼에 대한 상수 조건이어야 한다. 함수 변환(예: DATE_TRUNC(...))이 들어가면 프루닝이 안 될 수 있다.

#### 통계 갱신(ANALYZE)

CBO(비용 기반 옵티마이저)가 잘 동작하려면 테이블 통계가 최신이어야 한다.

-- 전체 테이블 통계 수집
ANALYZE iceberg.warehouse.orders;

-- 파티션 지정 (큰 테이블에서 증분 갱신)
ANALYZE iceberg.warehouse.orders
  WITH (partitions = ARRAY[ARRAY['2026-07-01']]);

SHOW STATS FOR orders로 현재 통계 상태를 확인한다. 행 수(row_count)와 null_fraction이 NULL이면 통계가 없다는 뜻이다.


Iceberg 커넥터 운영 주의점

상황원인조치
쿼리가 모든 파일을 읽음파티션 프루닝 실패EXPLAIN 출력에서 TableScanNode 확인, 파티션 컬럼 조건 점검
ANALYZE 후에도 CBO가 잘못된 계획통계가 오래됨SHOW STATS FOR 조회 후 재실행
파일 수가 폭발적으로 증가소량 삽입 반복Iceberg OPTIMIZE 명령으로 compaction 실행
메타데이터 파일 과다잦은 커밋expire_snapshots 프로시저로 오래된 스냅샷 정리
-- Compaction (소파일 병합)
ALTER TABLE events EXECUTE optimize(file_size_threshold => '128MB');

-- 오래된 스냅샷 삭제
ALTER TABLE events EXECUTE expire_snapshots(retention_threshold => '7d');

MySQL 커넥터

기본 설정

# etc/catalog/mysql_prod.properties
connector.name=mysql
connection-url=jdbc:mysql://mysql-host:3306
connection-user=trino_reader
connection-password=<password>
# 연결 풀 크기
connection-pool-min-size=1
connection-pool-max-size=30

Trino는 MySQL 드라이버를 통해 연결한다. 커넥터 이름을 mysql_prod, mysql_analytics처럼 구분해 여러 MySQL 클러스터에 동시 연결할 수 있다.

푸시다운(Pushdown)

Trino의 핵심 최적화는 연산을 데이터소스로 "밀어내리는(push down)" 것이다. 푸시다운이 될수록 Trino가 읽어야 할 데이터가 줄어든다.

Trino가 처리 (푸시다운 불가)
CASE/IF 일부 CHAR/VARCHAR 필터 복잡한 표현식 정규식 함수
조건부 푸시다운 (설정 필요)
집계 함수 (SUM/COUNT/AVG) JOIN 일부 (Connector Join) LIMIT/OFFSET
항상 푸시다운 (자동)
숫자/날짜 컬럼 WHERE 조건 컬럼 프로젝션 ORDER BY (일부)
MySQL 커넥터 푸시다운 계층

중요한 한계: Trino MySQL 커넥터는 CHAR/VARCHAR 컬럼에 대한 필터를 푸시다운하지 않는다. MySQL의 문자열 비교가 기본적으로 대소문자를 무시하는(case-insensitive) 방식이라 결과 정확성 보장이 안 되기 때문이다. 이를 모르고 WHERE name = 'alice'가 MySQL에서 걸릴 것이라 가정하면 전체 테이블 스캔이 Trino 단으로 올라온다.

집계 푸시다운 활성화

집계 푸시다운은 명시적으로 켜야 한다.

# etc/catalog/mysql_prod.properties
connector.name=mysql
connection-url=jdbc:mysql://mysql-host:3306
connection-user=trino_reader
connection-password=<password>
optimizer.aggregation-pushdown-enabled=true

혹은 세션 속성으로 임시 활성화:

SET SESSION mysql_prod.aggregation_pushdown_enabled = true;
SELECT customer_id, SUM(amount) FROM mysql_prod.sales.orders GROUP BY customer_id;

집계 푸시다운이 동작하면 EXPLAIN 결과에 Aggregate 노드가 사라지고 MySQL이 직접 SUM을 계산한다.

대형 IN 절 처리

WHERE id IN (1, 2, 3, ..., 10000) 같이 IN 절이 매우 큰 경우, Trino는 기본 256개를 초과하면 범위 조건(range predicate)으로 압축한다. 이 임계값을 조정할 수 있다.

# 임계값 높이기 (IN 절을 더 많이 푸시다운)
domain-compaction-threshold=1000

너무 큰 IN 절은 MySQL 쿼리 파싱 자체의 성능 문제를 일으킬 수 있으므로 균형이 필요하다.

운영 패턴

-- Trino를 통한 크로스-소스 조인: MySQL 운영 DB + Iceberg 분석 테이블
SELECT
    c.name,
    c.region,
    SUM(o.amount) AS total_amount
FROM mysql_prod.crm.customers c
JOIN iceberg.warehouse.orders o ON c.id = o.customer_id
WHERE o.order_date >= DATE '2026-01-01'
GROUP BY c.name, c.region
ORDER BY total_amount DESC
LIMIT 100;

이 쿼리에서 Trino는:

  1. MySQL customers를 읽을 때 컬럼 프로젝션(name, region, id만)을 푸시다운한다.
  2. Iceberg orders를 읽을 때 order_date 조건으로 파티션 프루닝을 적용한다.
  3. 테이블 크기 통계에 따라 broadcast join 또는 hash join을 선택한다.

Kafka 커넥터

Kafka 커넥터의 위치

Trino의 Kafka 커넥터는 실시간 이벤트 스트림을 배치 SQL로 탐색하는 용도에 맞다. "지금 이 토픽에 어떤 메시지가 있는지"를 빠르게 확인하거나, 스트림 데이터와 다른 소스를 조인할 때 유용하다.

주의: Kafka 커넥터는 스트리밍 처리 엔진이 아니다. 쿼리를 실행하는 순간의 오프셋 범위를 읽는 정적 스캔이다. 실시간 집계나 연속 처리는 Flink를 사용해야 한다.

기본 설정

# etc/catalog/kafka_prod.properties
connector.name=kafka
kafka.nodes=kafka-broker-1:9092,kafka-broker-2:9092
kafka.table-names=events.page_views,events.clicks
kafka.default-schema=events

kafka.table-names에 직접 토픽 이름을 나열하거나, kafka.table-description-dir에 JSON 스키마 파일을 두는 방식을 쓴다.

Confluent Schema Registry 연동

토픽 메시지가 Avro 형식이고 Confluent Schema Registry를 사용한다면:

connector.name=kafka
kafka.nodes=kafka-broker-1:9092
kafka.table-description-supplier=CONFLUENT
kafka.confluent-schema-registry-url=http://schema-registry:8081
kafka.confluent-schema-registry-client-cache-size=1000

kafka.table-description-supplier=CONFLUENT로 설정하면 Trino가 Schema Registry에서 최신 스키마를 자동으로 읽어 테이블 구조를 결정한다. 컬럼 이름과 타입이 자동으로 매핑된다.

제약사항: Confluent Schema Registry 모드에서는 Avro 형식만 지원되며, 삽입(INSERT)은 불가능하다. 읽기(SELECT) 전용이다.

내장 메타 컬럼

Kafka 커넥터는 Kafka 메시지의 메타 정보를 특수 컬럼으로 노출한다.

컬럼타입설명
_partition_idBIGINTKafka 파티션 번호
_partition_offsetBIGINT파티션 내 오프셋
_message_corruptBOOLEAN메시지 파싱 오류 여부
_keyVARCHAR메시지 키 (원시)
_key_corruptBOOLEAN키 파싱 오류
_timestampTIMESTAMP프로듀서 타임스탬프
_messageVARCHAR메시지 본문 (원시, 파싱 전)
-- 파티션별 최신 오프셋 확인
SELECT _partition_id, MAX(_partition_offset) AS latest_offset
FROM kafka_prod.events.page_views
GROUP BY _partition_id
ORDER BY _partition_id;

-- 특정 시간대 메시지 탐색
SELECT _partition_id, _partition_offset, user_id, url
FROM kafka_prod.events.page_views
WHERE _timestamp >= TIMESTAMP '2026-07-11 10:00:00'
  AND _timestamp <  TIMESTAMP '2026-07-11 11:00:00'
LIMIT 100;

Kafka 커넥터 아키텍처: 어떻게 읽는가

Coordinator 1. 파티션 목록 조회 2. 오프셋 범위 결정 3. Split 생성 후 배분 Split 배분 Worker 1 Split: partition 0 offset 0~9999 Worker 2 Split: partition 1 offset 0~12000 Kafka 클러스터 partition-0 (broker-1) partition-1 (broker-2) partition-2 (broker-3) Split = 파티션 하나의 특정 오프셋 범위 각 워커가 할당된 파티션의 오프셋 범위를 Consumer API로 직접 읽음
Trino Kafka 커넥터 읽기 흐름

Trino는 Kafka 파티션 하나를 여러 Split으로 나눌 수 있다. kafka.max-split-bytes(기본 64MB)로 Split 크기를 조절한다.


커넥터 비교: 운영 관점

Iceberg 읽기/쓰기 모두 가능 파티션 프루닝 강력 타임 트래블 지원 ANALYZE로 통계 갱신 필요 small files 관리 필수
MySQL 읽기/쓰기 가능 숫자·날짜 필터 푸시다운 집계 푸시다운(옵션) VARCHAR 필터 푸시다운 불가 운영 DB에 부하 주의
Kafka 읽기 전용(SR 모드) 오프셋 범위 기반 Split Schema Registry 연동 집계 푸시다운 없음 retention 지난 데이터 없음
세 커넥터 특성 요약

커넥터 조합 실전 패턴

패턴 1: 운영 DB + 분석 레이크 조인

-- MySQL 운영 DB의 최신 고객 정보 + Iceberg 이력 주문 데이터
SELECT
    c.tier,
    COUNT(DISTINCT o.customer_id) AS active_customers,
    AVG(o.total_amount) AS avg_order_value
FROM mysql_prod.crm.customers c
JOIN iceberg.warehouse.orders o ON c.id = o.customer_id
WHERE o.created_at >= DATE '2026-01-01'
GROUP BY c.tier;

customers는 수백만 행, orders는 수십억 행이라면 Trino는 작은 테이블(customers)을 broadcast join 대상으로 결정할 수 있다. 단, MySQL 커넥터에서 테이블 통계를 얻기 어려우면 CBO가 잘못된 판단을 할 수 있다. 이럴 때 힌트를 준다.

SELECT /*+ BROADCAST(c) */ ...

패턴 2: Kafka 이벤트 + MySQL 마스터 데이터 실시간 조회

-- 최근 1시간 에러 이벤트와 서비스 정보 조합
SELECT
    e.service_id,
    s.service_name,
    COUNT(*) AS error_count
FROM kafka_prod.events.errors e
JOIN mysql_prod.config.services s ON e.service_id = s.id
WHERE e._timestamp >= NOW() - INTERVAL '1' HOUR
  AND e.level = 'ERROR'
GROUP BY e.service_id, s.service_name
ORDER BY error_count DESC;

References

  • Trino Iceberg connector 공식 문서, https://trino.io/docs/current/connector/iceberg.html
  • Trino MySQL connector 공식 문서, https://trino.io/docs/current/connector/mysql.html
  • Trino Kafka connector 공식 문서, https://trino.io/docs/current/connector/kafka.html
  • Trino Pushdown 공식 문서, https://trino.io/docs/current/optimizer/pushdown.html
  • Trino Metastores 공식 문서, https://trino.io/docs/current/object-storage/metastores.html
  • Iceberg connector — Starburst Enterprise, https://docs.starburst.io/latest/connector/iceberg.html
  • Stackable Tech: Trino REST Catalog Discussion, https://github.com/orgs/stackabletech/discussions/41
  • Trino GitHub: Kafka Schema Registry support, https://github.com/trinodb/trino/issues/12195