Apache Beam·Flink·Hive QL·Presto로 데이터 처리 경로 설계하기
Apache Beam, Flink, Hive QL, Presto(Trino)의 처리 모델과 일관성, 카탈로그 운영 방식을 비교해 데이터 플랫폼의 역할을 설계한다.
2026-08-14 · 최초 발행 2025-10-14
하나의 엔진으로 묶기 어려운 데이터 처리 경로
실시간 적재, 배치 변환, 웨어하우스 갱신, 대화형 분석은 같은 데이터 레이크를 바라보더라도 요구하는 실행 방식이 다르다. Apache Beam은 파이프라인을 기술하는 계층을, Apache Flink는 상태를 가진 스트리밍 처리를, Hive QL은 배치 SQL 변환을, Presto(Trino)는 여러 저장소를 넘나드는 대화형 질의를 맡는다.
이 네 도구를 함께 쓸 때는 기능 목록보다 데이터가 어느 시점에 어떤 일관성을 가져야 하는지부터 정해야 한다.
Apache Beam: 파이프라인과 실행 엔진을 분리한다
Apache Beam은 데이터 처리 파이프라인을 모델, SDK, 런너로 나눈 통합 프로그래밍 모델이다. 동일한 파이프라인을 여러 실행 엔진에서 돌릴 수 있어 코드 이식성과 표준화에 유리하다.
배치와 스트리밍을 하나의 모델 안에서 다루며, 이벤트 타임, 윈도잉, 트리거, 정합성 semantics를 제공한다. 실제 실행은 Flink, Spark, Google Cloud Dataflow 등의 런너가 담당한다.
Apache Flink: 상태를 유지하는 스트리밍 경로
Flink는 스트리밍 우선 엔진으로, 정확히 한 번 처리와 대규모 상태 관리, 저지연 처리에 초점이 맞춰져 있다. 체크포인팅, 세이브포인트, State Backend(rocksdb 등)를 기반으로 고가용성과 회복성을 구성한다.
DataStream API와 Flink SQL/Table API를 모두 지원하므로 실시간 분석과 CEP에 사용할 수 있다. 특히 상태를 보존한 채 스트림을 처리하고 복구해야 하는 흐름에 맞는다.
Hive QL: 배치 변환과 웨어하우스 갱신
Hive QL은 Hadoop 에코시스템을 위한 SQL 계열 언어다. Hive Metastore를 중심으로 스키마 온 리드 방식을 제공하며, Tez, MapReduce, Spark에서 실행할 수 있다.
ORC와 Parquet, 파티셔닝, 버킷팅 최적화를 활용할 수 있고, ORC 기반 ACID 테이블에서는 트랜잭션 업데이트도 가능하다. 대규모 배치 ETL과 웨어하우스 변환 작업에 적합하다.
Presto(Trino): 저장소를 가로지르는 대화형 SQL
Presto(Trino)는 분산 MPP SQL 쿼리 엔진이다. 데이터 레이크와 RDBMS를 다양한 커넥터로 연결해 조인하거나 연합 질의할 수 있다. 메모리 기반 파이프라인 처리로 대화형 쿼리 성능을 확보하며, 레이크하우스와 라이트 아키텍처에 어울린다.
PrestoSQL이 Trino로 바뀌고 PrestoDB가 병행되는 명칭 분화 이력이 있으므로, 배포판과 커넥터의 호환성을 확인할 필요가 있다.
처리 모델과 메타데이터를 함께 설계해야 한다
Beam은 파이프라인 추상화와 런너 분리를 통해 코드 재사용성과 벤더 종속 최소화에 집중한다. Flink는 네이티브 런타임 최적화와 정확히 한 번 보장을 제공하는 상태 기반 스트리밍 엔진이다. Hive와 Presto는 SQL 중심이지만, Hive는 배치 ETL에, Presto는 대화형·연합 분석에 더 맞는다.
Hive Metastore나 Glue Catalog를 중심으로 스키마를 관리하면 Presto와 Flink SQL에서도 이를 활용할 수 있다. ORC, Parquet, Avro로 포맷을 표준화하고 파티셔닝과 Z-order/Clustering을 적용하면 프루닝을 최적화할 수 있다. 다만 커넥터와 카탈로그를 이중화하면 메타데이터 일관성을 관리하는 체계가 필요하다.
Flink는 체크포인트와 트랜잭셔널 싱크로 정확히 한 번 출력을 보장한다. Hive는 ORC ACID 테이블의 스냅샷과 락 관리로 배치 머지와 업데이트를 지원한다. Presto는 주로 읽기와 분석 경로를 최적화하며, 쓰기 일관성은 아이스버그나 델타 등의 테이블 포맷 트랜잭션에 의존한다.
운영 환경에서는 YARN 또는 K8s로 배치와 스트리밍 워크로드를 격리하고, 메트릭·로그·트레이스를 Grafana/Prometheus와 OpenTelemetry로 일원화할 수 있다. 파일 컴팩션, 파티션 관리, 쿼리 캐시, 코스트 기반 옵티마이저도 비용과 성능을 조정하는 수단이다.
데이터 레이크 위에서 역할을 나누는 방식
실시간 ETL과 집계에서는 Flink가 Kafka에서 Lake로 적재하고 윈도 집계를 수행하며 정확히 한 번 처리를 담당할 수 있다. 공통 로직은 Beam 파이프라인으로 정의하고, 개발과 테스트를 거친 뒤 Flink 또는 Dataflow에서 실행한다. 최신 집계 결과는 Presto로 조회해 대시보드와 SLO 경로에 제공한다.
대규모 배치 ETL, CTAS, 머지는 Hive QL이 맡고 ORC 또는 Parquet 포맷을 정착시킬 수 있다. 임시 분석과 페더레이티드 조인은 Presto로 제공해 데이터 마트를 대체하거나 보완한다. 이때 메타스토어를 단일화하면 거버넌스와 스키마 에볼루션을 통제하기 쉽다.
머신러닝 피처 파이프라인에서는 Beam 또는 Flink가 실시간 피처 스트림 변환과 서빙 준비를 맡고, Hive는 학습용 스냅샷 생성과 데이터 품질 검증 룰 적용을 담당한다. Presto는 피처 탐색, 샘플링, 오프라인 검증에 사용할 수 있다.
선택 기준을 한 표에서 비교하기
| 도구 | 처리 모델 | 성능(지연/처리량) | 확장성 | 일관성 | 안정성 | 운영 편의 |
|---|---|---|---|---|---|---|
| Apache Beam | 배치+스트림(추상화) | 런너에 의존, 코드 이식성 우수 | 런너 확장성 활용 | 이벤트 타임·트리거 모델 | 런너의 체크포인트/재시도에 의존 | 단일 코드베이스·멀티 런타임 |
| Apache Flink | 스트림 우선+배치 | 저지연, 고처리량, 상태 관리 강점 | K8s/YARN 수평 확장 | 정확히 한 번 보장 | 체크포인트/세이브포인트/재처리 | SQL·API 혼용, 세밀한 튜닝 |
| Hive QL | 배치 | 고처리량, 높은 지연 | HDFS/오브젝트 스토리지 기반 | ACID(ORC) 선택적 적용 | 실패 시 재시도·맵리듀스/Tez 복원 | ETL 일관성, 운영 패턴 성숙 |
| Presto(Trino) | 대화형 쿼리 | 초저지연(초~수십초), 메모리 최적화 | 워커 수평 확장 | 읽기 일관성, 쓰기는 포맷 의존 | 쿼리 재시도·코디네이터 HA 구성 | 카탈로그/커넥터 기반 운영 용이 |
저장·처리·제공 경로
파이프라인과 SQL 예시
Beam에서는 이벤트 타임 윈도 집계와 정확히 한 번 싱크 구성을 다룬다. 전제는 Python 3.10+, apache-beam[gcp]==2.53+, Flink 1.17+, Kafka 커넥터 구성이다.
# pip install apache-beam==2.53.0 apache-beam[interactive]
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
opts = PipelineOptions([
"--runner=FlinkRunner",
"--flink_version=1.17",
"--flink_master=localhost:8081",
"--parallelism=4",
])
with beam.Pipeline(options=opts) as p:
(p
| "ReadJSON" >> beam.io.ReadFromKafka(
consumer_config={"bootstrap.servers":"localhost:9092"},
topics=["events"])
| "Parse" >> beam.Map(lambda kv: json.loads(kv[1]))
| "WithTS" >> beam.Map(lambda e: beam.window.TimestampedValue(e, e["event_time"]))
| "Window" >> beam.WindowInto(beam.window.FixedWindows(60))
| "KeyBy" >> beam.Map(lambda e: (e["type"], 1))
| "Count" >> beam.CombinePerKey(sum)
| "ToParquet" >> beam.io.parquetio.WriteToParquet(
file_path_prefix="s3://lake/agg/min1",
schema={"fields":[{"name":"type","type":"BYTE_ARRAY"},
{"name":"cnt","type":"INT64"}]}))
Flink SQL에서는 Kafka/Filesystem 커넥터와 Table SQL Client를 전제로 워터마크와 텀블링 윈도우를 사용한다. Flink 1.17+ 환경을 기준으로 한다.
CREATE TABLE events (
user_id STRING,
etime TIMESTAMP(3),
WATERMARK FOR etime AS etime - INTERVAL '5' SECOND
) WITH (...);
CREATE TABLE agg_min1 (
win_start TIMESTAMP(3),
cnt BIGINT
) WITH (...);
INSERT INTO agg_min1
SELECT
TUMBLE_START(etime, INTERVAL '1' MINUTE) AS win_start,
COUNT(*) AS cnt
FROM events
GROUP BY TUMBLE(etime, INTERVAL '1' MINUTE);
Hive QL의 배치 경로에서는 Hive 3.x와 ORC ACID 활성화(transactional=true)를 전제로 파티션 CTAS와 MERGE를 적용한다.
SET hive.txn.manager=org.apache.hadoop.hive.ql.lockmgr.DbTxnManager;
SET hive.compactor.initiator.on=true;
CREATE TABLE fact_daily
STORED AS ORC
TBLPROPERTIES ('transactional'='true')
AS
SELECT * FROM staging_fact WHERE dt='${proc_date}';
MERGE INTO fact_daily t
USING delta d
ON t.id=d.id AND t.dt=d.dt
WHEN MATCHED THEN UPDATE SET amount=d.amount
WHEN NOT MATCHED THEN INSERT VALUES(d.id, d.amount, d.dt);
Presto 또는 Trino에서는 Hive/Glue 카탈로그 연결을 전제로 레이크 테이블 조인과 근사 집계를 수행할 수 있다.
SELECT d.region, approx_distinct(f.user_id) AS uu
FROM hive.lake.fact_daily f
JOIN hive.ref.dim_region d ON f.region_id = d.id
WHERE f.dt BETWEEN date '2025-09-01' AND date '2025-09-30'
GROUP BY d.region
ORDER BY uu DESC
LIMIT 20;
운영 지표로 확인할 효과
스트리밍 집계는 p95 지연 15초 달성, 대화형 쿼리는 115초 응답 범위를 확보하는 것이 목표가 될 수 있다. 스토리지 포맷과 파티셔닝을 최적화하면 스캔 바이트를 3070% 절감하고 쿼리 시간을 2050% 단축할 수 있다.
정확히 한 번 파이프라인과 ACID 테이블을 조합하면 데이터 오류와 중복을 90% 이상 줄일 수 있다. 파일 포맷과 카탈로그를 먼저 표준화한 뒤 스트리밍 파이프라인을 구축하고, 대화형 쿼리를 확장하는 순서로 도입하면 각 도구의 역할을 분리하기 쉽다.