스트림과 테이블 사이의 경계
Apache Flink의 통합 테이블·스트림 모델은 강력하지만, 그 경계를 개발자가 직접 제어하기 어렵다는 한계가 있었다. 동적 테이블(dynamic table)과 체인지로그 스트림(changelog stream) 사이의 변환이 프레임워크 내부에서 암묵적으로 일어났기 때문이다. CDC 소스를 받아 커스텀 형식의 업서트 스트림을 만들거나, 역방향으로 특정 변경 스트림을 다시 테이블로 소비하는 작업은 DataStream API를 섞어야만 가능했다.
2026년 6월 25일 릴리스된 Apache Flink 2.3.0은 이 경계를 SQL로 명시적으로 다룰 수 있게 한다. FROM_CHANGELOG와 TO_CHANGELOG 두 연산자, 그리고 ALTER MATERIALIZED TABLE에 추가된 풀 DDL 지원이 핵심이다. 아울러 Hadoop 의존성을 완전히 제거한 네이티브 S3 플러그인(flink-s3-fs-native)이 실험적으로 제공된다.
FROM_CHANGELOG: 체인지로그 스트림을 동적 테이블로
-- CDC 형식의 append-only 스트림이 들어오는 경우
-- op_col: 각 행의 변경 유형을 담은 컬럼 ('+I', '-U', '+U', '-D' 등)
SELECT *
FROM FROM_CHANGELOG(
TABLE cdc_source_stream,
op_col => 'op',
op_mapping => MAP['+I', INSERT, '-U', UPDATE_BEFORE, '+U', UPDATE_AFTER, '-D', DELETE]
);FROM_CHANGELOG는 append-only 형태로 들어오는 스트림을 Flink의 동적 테이블로 변환한다. 스트림의 각 행에는 변경 유형을 나타내는 컬럼(op_col)이 있고, op_mapping으로 이 값을 Flink 내부의 RowKind(INSERT, UPDATE_BEFORE, UPDATE_AFTER, DELETE)에 매핑한다.
이전까지는 이런 변환을 하려면 DataStream API에서 커스텀 MapFunction을 작성하거나 특수 커넥터를 구현해야 했다. FROM_CHANGELOG 하나로 Flink SQL 안에서 처리할 수 있게 됐다.
op_mapping에 정의되지 않은 operation 코드를 가진 행을 어떻게 처리할지도 제어 가능하다. on_missing_operation_type 파라미터로 IGNORE, ERROR, 또는 기본값으로 처리하도록 설정한다.
TO_CHANGELOG: 동적 테이블을 append-only 스트림으로
TO_CHANGELOG는 반대 방향이다. retract 스트림이나 upsert 스트림 형태의 동적 테이블을 append-only 형식의 체인지로그로 변환한다.
-- 집계 결과를 append-only 아카이브 싱크에 쓰기
INSERT INTO audit_log_sink
SELECT *
FROM TO_CHANGELOG(
TABLE (SELECT region, COUNT(*) AS cnt FROM orders GROUP BY region),
op_col => 'change_type'
);TO_CHANGELOG가 중요한 이유는 append-only 싱크와의 호환성이다. Kafka에 데이터를 아카이브하거나 객체 스토리지에 감사 로그(audit log)를 쓸 때, 소스가 retract 스트림이면 직접 연결이 불가능했다. TO_CHANGELOG를 거치면 변경 유형 컬럼을 포함한 일반 append 행으로 직렬화된다.
두 연산자의 대칭 구조
두 연산자는 대칭 관계다. FROM_CHANGELOG는 외부 CDC 포맷을 Flink의 테이블 모델로 흡수하고, TO_CHANGELOG는 Flink 내부의 동적 테이블을 append 친화적인 외부 싱크로 내보낸다. SQL 파이프라인 양 끝단에서 체인지로그 형식을 명시적으로 제어할 수 있게 됐다.
ALTER MATERIALIZED TABLE: 풀 DDL 지원
Flink 2.3.0은 Materialized Table의 DDL 능력을 일반 테이블 수준으로 끌어올렸다.
컬럼 변경
-- 컬럼 추가
ALTER MATERIALIZED TABLE order_summary
ADD COLUMN avg_amount DOUBLE COMMENT '평균 주문금액';
-- 컬럼 수정
ALTER MATERIALIZED TABLE order_summary
MODIFY COLUMN avg_amount DECIMAL(10,2);
-- 컬럼 삭제
ALTER MATERIALIZED TABLE order_summary
DROP COLUMN deprecated_col;
-- 테이블 이름 변경
ALTER MATERIALIZED TABLE order_summary RENAME TO order_metrics;이전 버전에서는 Materialized Table의 스키마를 바꾸려면 삭제 후 재생성해야 했다. 2.3.0부터는 일반 ALTER TABLE과 동일한 워크플로로 점진적으로 발전시킬 수 있다.
START_MODE: 재처리 경계 제어
쿼리가 변경되어 Materialized Table을 다시 refresh해야 할 때, 어디서부터 다시 읽을지를 START_MODE로 명시한다.
ALTER MATERIALIZED TABLE order_metrics
RESUME JOB WITH (
'start_mode' = 'earliest' -- 소스 최초 오프셋부터 재처리
-- 또는 'start_mode' = 'latest' -- 현재 오프셋부터
-- 또는 'start_mode' = 'specific_offset' -- 지정 오프셋
);특히 유용한 옵션은 이전 잡이 멈춘 정확한 소스 오프셋부터 재시작하는 것이다. 이전 refresh 잡 인스턴스의 상태를 보존해 두었다가, 쿼리 변경 후에도 데이터 중복이나 누락 없이 이어서 처리할 수 있다.
flink-s3-fs-native: Hadoop 없는 S3
Flink의 기존 S3 플러그인(flink-s3-fs-hadoop, flink-s3-fs-presto)은 모두 Hadoop 파일시스템 레이어나 Presto 코드에 의존했다. 이는 불필요한 JAR 의존성을 끌어오고, 비동기 I/O 지원이 제한적이었다.
Flink 2.3.0은 flink-s3-fs-native를 실험적 기능으로 도입했다.
| 항목 | 기존 플러그인 | flink-s3-fs-native |
|---|---|---|
| 기반 SDK | AWS SDK v1 + Hadoop | AWS SDK v2 (직접) |
| I/O 모델 | 동기 | 비동기 (논블로킹) |
| Hadoop 의존성 | 필요 | 없음 |
| IAM Roles for Service Accounts | 제한적 | 완전 지원 |
| RecoverableWriter | 별도 구현 | 통합 구현 |
| 상태 | 안정 | 실험적 (2.3.0 기준) |
RecoverableWriter(정확히-한-번 쓰기를 위한 인터페이스)가 FileSystem과 동일한 구현 안에 통합된 점도 변화다. 기존에는 두 인터페이스를 따로 다루는 구현이 필요했다.
<!-- flink-conf.yaml 또는 플러그인 디렉토리 설정 -->
<!-- plugins/flink-s3-fs-native/ 아래에 JAR 배치 -->
<!-- s3.endpoint, s3.path-style-access 등 기존 설정 그대로 사용 가능 -->2.3.0 기준 실험적 단계이므로 프로덕션 워크로드에 적용하기 전 충분한 스테이징 검증이 필요하다. AWS SDK v2 기반의 자격 증명 프로바이더(IRSA, AssumeRole 등)와의 호환성 검증도 포함해야 한다.
업그레이드 시 주의사항
레거시 객체 스토리지 어댑터 제거 (Trino 방향과 유사)
Flink 2.3.0은 오래된 레거시 객체 스토리지 어댑터를 계속 정비하고 있다. 사용 중인 플러그인이 deprecated 경고를 내는지 확인한다. 레거시 어댑터에서 네이티브 플러그인으로의 마이그레이션이 필요할 수 있다.
Materialized Table 기존 잡과의 호환성
ALTER MATERIALIZED TABLE의 새 DDL 기능은 이미 실행 중인 refresh 잡에 영향을 준다. 컬럼을 변경하기 전에 실행 중인 잡을 중단하고, 변경 후 RESUME JOB으로 재시작하는 순서를 지킨다.
FROM_CHANGELOG op_mapping 설계
op_mapping을 잘못 정의하면 행이 잘못된 RowKind로 처리된다. 특히 UPDATE_BEFORE와 UPDATE_AFTER를 모두 매핑해야 하는 경우, 소스 CDC 포맷이 두 행을 쌍으로 보내는지 단일 행으로 보내는지 확인한다.
운영 전환 체크리스트
- [ ] Flink 버전 확인: 2.3.0은 Java 11·17 지원, Java 8 지원 종료
- [ ] CDC 파이프라인:
FROM_CHANGELOGop_mapping을 소스 포맷(Debezium, Canal, Maxwell)에 맞게 설계 - [ ] audit/archive 파이프라인: retract 소스 →
TO_CHANGELOG→ Kafka/S3 append 싱크 경로 검증 - [ ] Materialized Table 컬럼 변경 전: 실행 중인 refresh 잡 중단 (
STOP JOB WITH SAVEPOINT) - [ ]
ALTER MATERIALIZED TABLE ... RESUME JOB후START_MODE선택 기준 문서화 - [ ] S3 플러그인 교체 고려:
flink-s3-fs-native스테이징 테스트 (프로덕션 전 충분한 검증) - [ ] IRSA(IAM Roles for Service Accounts) 환경:
flink-s3-fs-native의 AWS SDK v2 자격증명 체인 검증
References
- https://flink.apache.org/2026/06/25/apache-flink-2.3.0-release-announcement/
- https://nightlies.apache.org/flink/flink-docs-master/docs/sql/materialized-table/statements/
- https://nightlies.apache.org/flink/flink-docs-master/release-notes/
- https://issues.apache.org/jira/browse/FLINK-30975
- https://developer.confluent.io/courses/flink-sql/changelog-processing/