Apache Flink 2.3: FROM_CHANGELOG·TO_CHANGELOG 연산자와 Materialized Table DDL 완성으로 스트리밍-배치 경계를 좁힌 방법
# Apache Flink 2.3: FROM_CHANGELOG·TO_CHANGELOG 연산자와 Materialized Table DDL 완성으로 스트리밍-배치 경계를 좁힌 방법
요약
Apache Flink 2.3.0이 2026년 6월 25일 출시됐다. 15개의 FLIP을 구현했으며 핵심 변경은 세 가지다. 첫째, FROM_CHANGELOG·TO_CHANGELOG SQL 연산자가 추가돼 CDC 스트림과 일반 스트림을 SQL 레벨에서 변환할 수 있게 됐다. 둘째, Materialized Table에 컬럼 명시 정의와 ALTER DDL이 붙어 테이블 생성 이후의 스키마 관리가 가능해졌다. 셋째, AWS SDK v2 기반 네이티브 S3 파일시스템이 실험적으로 도입돼 hadoop-aws 의존성 없이 고성능 오브젝트 스토리지 접근이 가능해졌다. 스트리밍과 배치를 하나의 SQL로 통합하려는 흐름이 이번 릴리스에서 한 단계 더 진전됐다.
배경
Flink 2.0(2025년 초)은 Unified Runtime을 확립했고, 2.2(2025년 12월)는 ML_PREDICT·VECTOR_SEARCH 같은 AI 통합 SQL 함수와 Delta Join 개선을 가져왔다. 2.3은 그 위에서 파이프라인 구성 단계의 문제를 다룬다.
실무에서 CDC 소스(MySQL binlog, Debezium 등)를 Flink로 받으면 결과가 changelog 스트림이 된다. 이 스트림을 다시 집계하거나 외부 싱크에 쓰기 전에 일반 append 스트림으로 변환해야 하는 경우가 많지만, 기존에는 별도 UDF나 DataStream API를 써야 했다. FROM_CHANGELOG·TO_CHANGELOG는 이 변환을 SQL 안으로 가져온 것이다.
Materialized Table은 Flink 2.1에서 도입된 기능으로 "이 SQL 결과를 주기적으로 갱신해 저장하라"는 의미다. 그런데 초기 구현에선 CREATE할 때 컬럼 타입이나 워터마크를 직접 지정할 수 없었고, 한 번 만든 뒤에는 DDL로 변경이 불가능했다. 2.3에서 이 두 가지가 모두 해결됐다.
FROM_CHANGELOG · TO_CHANGELOG 연산자
문제
CDC 소스를 Flink로 읽으면 각 레코드는 삽입(+I), 수정 전(-U), 수정 후(+U), 삭제(-D) 중 하나의 RowKind를 가진다. 이 changelog 스트림을 그대로 집계하거나 window 연산에 넣으면 예상치 못한 결과가 나올 수 있다. 반대로, 일반 append 스트림 결과를 downstream 테이블에 CDC 형식으로 전달하고 싶을 때도 변환이 필요했다.
FROM_CHANGELOG
SELECT *
FROM FROM_CHANGELOG(
TABLE orders, -- changelog 스트림 소스
'retract' -- 변환 전략: 'retract' 또는 'upsert'
)FROM_CHANGELOG는 changelog 스트림을 일반 append 스트림으로 변환한다. 전략은 두 가지다.
retract:-U(before) 레코드는 삭제 이벤트로,+U(after)는 삽입 이벤트로 분리한다.upsert: primary key 기준으로-D·-U를 absorb해 순수 upsert 스트림을 만든다.
TO_CHANGELOG
INSERT INTO target_changelog_table
SELECT *
FROM TO_CHANGELOG(
TABLE aggregation_result,
PRIMARY KEY (order_id) NOT ENFORCED
)TO_CHANGELOG는 반대 방향이다. 일반 append 결과를 downstream이 기대하는 changelog 형식으로 wrapping한다. primary key를 지정하면 upsert 시멘틱으로, 지정하지 않으면 insert-only 스트림으로 처리된다.
파이프라인 합성 예시
-- MySQL CDC → 변환 → 집계 → 아웃바운드 CDC 싱크
CREATE VIEW order_events AS
SELECT order_id, product_id, quantity
FROM FROM_CHANGELOG(TABLE mysql_orders_cdc, 'upsert');
CREATE MATERIALIZED TABLE daily_product_qty
REFRESH EVERY 1 HOUR AS
SELECT product_id, SUM(quantity) AS total_qty, CURRENT_DATE AS dt
FROM order_events
GROUP BY product_id, CURRENT_DATE;이렇게 하면 CDC 소스에서 집계 Materialized Table로 이어지는 파이프라인을 DataStream API 없이 순수 SQL로 구성할 수 있다.
Materialized Table DDL 완성
CREATE 시 컬럼 명시
이전 버전에선 CREATE MATERIALIZED TABLE ... AS SELECT ...만 가능해서 결과 스키마는 SELECT가 결정했다. 2.3에서는 일반 CREATE TABLE처럼 컬럼을 직접 정의할 수 있다.
CREATE MATERIALIZED TABLE hourly_sales (
product_id BIGINT,
total_revenue DECIMAL(18, 2),
window_start TIMESTAMP(3),
window_end TIMESTAMP(3),
WATERMARK FOR window_start AS window_start - INTERVAL '5' SECOND,
PRIMARY KEY (product_id, window_start) NOT ENFORCED
)
REFRESH EVERY 1 HOUR
AS
SELECT
product_id,
SUM(price * quantity) AS total_revenue,
TUMBLE_START(...) AS window_start,
TUMBLE_END(...) AS window_end
FROM transactions
GROUP BY product_id, TUMBLE(proc_time, INTERVAL '1' HOUR);워터마크와 primary key를 명시하면 downstream 엔진이 late data 처리와 upsert 병합을 직접 이해할 수 있다.
ALTER MATERIALIZED TABLE
-- 컬럼 추가
ALTER MATERIALIZED TABLE hourly_sales ADD COLUMN region STRING;
-- 컬럼 타입 변경
ALTER MATERIALIZED TABLE hourly_sales MODIFY COLUMN total_revenue DECIMAL(24, 4);
-- 컬럼 삭제
ALTER MATERIALIZED TABLE hourly_sales DROP COLUMN region;
-- 이름 변경
ALTER MATERIALIZED TABLE hourly_sales RENAME TO hourly_revenue_by_product;이전에는 스키마를 변경하려면 테이블을 드롭하고 재생성해야 했다. ALTER가 생기면서 운영 중인 Materialized Table의 스키마를 무중단으로(혹은 최소 중단으로) 조정할 수 있게 됐다.
세밀한 갱신 전략
REFRESH 절에 새로운 옵션이 추가됐다.
-- 전체 재계산 대신 신규 파티션만 계산
CREATE MATERIALIZED TABLE monthly_summary
REFRESH EVERY 1 DAY PARTITION BY dt
AS SELECT ...;
-- 수동 트리거
ALTER MATERIALIZED TABLE monthly_summary REFRESH PARTITION dt = '2026-08-01';이 파티션 단위 갱신은 배치 성 워크로드에서 불필요한 전체 재계산을 피하게 해준다.
네이티브 S3 파일시스템 (실험적)
기존 문제
기존 Flink의 S3 접근은 hadoop-aws JAR에 의존했다. 이 패키지는 Hadoop의 버전 의존성을 끌고 들어와서, Flink만 단독 배포하고 싶어도 거대한 Hadoop 생태계 의존성을 함께 패키징해야 했다. 또한 AWS SDK v1 기반이라 성능과 API 호환성에서 AWS SDK v2와 격차가 있었다.
flink-s3-fs-aws
2.3부터 flink-s3-fs-aws 플러그인이 도입됐다. AWS SDK v2를 직접 사용하며 hadoop-aws 없이 S3 읽기/쓰기가 가능하다. 설정 방식은 기존 flink-s3-fs-hadoop과 동일하게 유지해 마이그레이션 비용을 낮췄다.
# flink-conf.yaml
fs.s3a.endpoint: s3.ap-northeast-2.amazonaws.com
fs.s3a.aws.credentials.provider: com.amazonaws.auth.InstanceProfileCredentialsProvider실험적 상태이므로 프로덕션 배포 전에 충분한 검증이 필요하지만, Flink를 Hadoop 없이 경량 배포하려는 환경(컨테이너, 서버리스)에서 이점이 크다.
Adaptive Partition Selection
분산 스트리밍 환경에서 upstream 태스크의 처리 속도가 고르지 않을 때, 특정 파티션에 데이터가 몰려 backpressure가 발생한다. 기존에는 파티션 할당이 고정돼 있어 느린 downstream 파티션이 전체 파이프라인 속도를 끌어내렸다.
2.3의 Adaptive Partition Selection은 런타임에 downstream 파티션의 부하를 모니터링해 덜 바쁜 파티션으로 데이터를 유도한다. 특히 셔플(shuffle) 단계가 있는 배치성 쿼리와 긴 window 집계에서 효과가 크다.
활성화 방법:
table.optimizer.adaptive-partition-selection.enabled: true운영자·엔지니어가 알아야 할 것
FROM_CHANGELOG 전략 선택: CDC 소스에 DELETE 이벤트가 있다면 retract를, upsert 시멘틱으로 downstream에 전달할 때는 upsert를 선택한다. 잘못 선택하면 downstream 집계에서 이중 계산이 발생한다.
Materialized Table ALTER와 갱신 일관성: ALTER MATERIALIZED TABLE로 컬럼을 추가하면 기존에 저장된 파티션에는 해당 컬럼이 없다. 이 불일치를 허용하는지 downstream 쿼리 엔진을 먼저 확인해야 한다. Iceberg나 Delta Lake처럼 스키마 진화를 지원하는 포맷을 사용하면 구 파티션도 null 값으로 읽힌다.
네이티브 S3 플러그인 실험적 상태: flink-s3-fs-aws는 아직 실험 단계다. 기존 flink-s3-fs-hadoop과 병행 사용이 안 되므로, 팟(Pod) 환경에서 플러그인을 교체할 때 plugins/ 디렉터리에 둘 다 두지 않도록 주의한다.
Adaptive Partition Selection과 정렬 보장: 이 기능은 처리량 최적화에 초점이 있고, 레코드 순서 보장을 하지 않는다. 이벤트 시간 정렬이 필요한 파이프라인에서는 활성화 전에 영향을 평가해야 한다.
버전 호환성: Flink 2.3은 Kafka Connector 3.4, JDBC Connector 3.3, Hive 3.1.x를 공식 지원한다. 기존 2.2 기반 Flink 애플리케이션을 이전할 때 FROM_CHANGELOG·TO_CHANGELOG 관련 카탈로그 메타데이터를 초기화해야 할 수 있다.
References
- Apache Flink 2.3.0 Release Announcement (2026-06-25): https://flink.apache.org/2026/06/25/apache-flink-2.3.0-release-announcement/
- Flink 2.3 Release Notes: https://nightlies.apache.org/flink/flink-docs-release-2.3/release-notes/flink-2.3/
- FLIP-468: FROM_CHANGELOG / TO_CHANGELOG operators
- FLIP-455: Materialized Table explicit column definitions and ALTER DDL
- Apache Flink GitHub Releases: https://github.com/apache/flink/releases/tag/release-2.3.0