Trino 커넥터 심화: Iceberg, MySQL, Kafka 커넥터 운영
커넥터는 Trino의 데이터 접속 계층이다
1편에서 커넥터 프레임워크의 구조를 살펴봤다. 이번 편에서는 실제 데이터 플랫폼에서 가장 자주 마주치는 세 커넥터를 깊이 다룬다.
- Iceberg 커넥터: S3나 HDFS 위에 올린 Lakehouse 레이어. 스냅샷, 파티션 프루닝, 스키마 발전(evolution).
- MySQL 커넥터: 운영 DB에서 직접 데이터를 읽거나 Trino를 통한 크로스-소스 조인.
- Kafka 커넥터: 실시간 토픽을 SQL로 조회. 스트림-배치 브리징.
각 커넥터의 설정, 한계, 운영 주의점을 구체적으로 짚는다.
Iceberg 커넥터
카탈로그 타입 선택
Iceberg는 테이블 메타데이터를 별도의 카탈로그에 저장한다. Trino의 iceberg.catalog.type 속성으로 어떤 카탈로그를 쓸지 지정한다.
HMS 카탈로그 설정
# etc/catalog/iceberg.properties
connector.name=iceberg
iceberg.catalog.type=hive_metastore
hive.metastore.uri=thrift://hms-host:9083
hive.metastore.authentication.type=NONEhive.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=30Trino는 MySQL 드라이버를 통해 연결한다. 커넥터 이름을 mysql_prod, mysql_analytics처럼 구분해 여러 MySQL 클러스터에 동시 연결할 수 있다.
푸시다운(Pushdown)
Trino의 핵심 최적화는 연산을 데이터소스로 "밀어내리는(push down)" 것이다. 푸시다운이 될수록 Trino가 읽어야 할 데이터가 줄어든다.
중요한 한계: 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는:
- MySQL
customers를 읽을 때 컬럼 프로젝션(name, region, id만)을 푸시다운한다. - Iceberg
orders를 읽을 때order_date조건으로 파티션 프루닝을 적용한다. - 테이블 크기 통계에 따라 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=eventskafka.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=1000kafka.table-description-supplier=CONFLUENT로 설정하면 Trino가 Schema Registry에서 최신 스키마를 자동으로 읽어 테이블 구조를 결정한다. 컬럼 이름과 타입이 자동으로 매핑된다.
제약사항: Confluent Schema Registry 모드에서는 Avro 형식만 지원되며, 삽입(INSERT)은 불가능하다. 읽기(SELECT) 전용이다.
내장 메타 컬럼
Kafka 커넥터는 Kafka 메시지의 메타 정보를 특수 컬럼으로 노출한다.
| 컬럼 | 타입 | 설명 |
|---|---|---|
_partition_id | BIGINT | Kafka 파티션 번호 |
_partition_offset | BIGINT | 파티션 내 오프셋 |
_message_corrupt | BOOLEAN | 메시지 파싱 오류 여부 |
_key | VARCHAR | 메시지 키 (원시) |
_key_corrupt | BOOLEAN | 키 파싱 오류 |
_timestamp | TIMESTAMP | 프로듀서 타임스탬프 |
_message | VARCHAR | 메시지 본문 (원시, 파싱 전) |
-- 파티션별 최신 오프셋 확인
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 커넥터 아키텍처: 어떻게 읽는가
Trino는 Kafka 파티션 하나를 여러 Split으로 나눌 수 있다. kafka.max-split-bytes(기본 64MB)로 Split 크기를 조절한다.
커넥터 비교: 운영 관점
커넥터 조합 실전 패턴
패턴 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