플링크 테이블 유지보수
플링크 테이블 유지보수 (Flink TableMaintenance)
이 문서에서는 Flink 환경에서 아이스버그 테이블 유지보수 작업을 수행하는 방법을 알려드릴게요. 배치 모드의 rewrite files 작업부터, 스트리밍 모드에서 스냅샷 만료·데이터 파일 컴팩션·고아 파일 정리를 자동화하는 TableMaintenance API, 그리고 커밋 후(post-commit) 유지보수 통합까지 코드 예시와 구성 옵션을 폭넓게 살펴볼게요.
출처: 문서
본문
Flink 테이블 유지보수 배치 모드 (Flink Table Maintenance BatchMode)
Rewrite files 작업 (Rewrite files action)
아이스버그는 Flink 배치 작업을 제출해서 작은 파일을 큰 파일로 재작성하는 API를 제공해요. 이 Flink 작업의 동작은 스파크의 rewriteDataFiles와 같아요.
import org.apache.iceberg.flink.actions.Actions;
TableLoader tableLoader = TableLoader.fromCatalog(
CatalogLoader.hive("my_catalog", configuration, properties),
TableIdentifier.of("database", "table")
);
Table table = tableLoader.loadTable();
RewriteDataFilesActionResult result = Actions.forTable(table)
.rewriteDataFiles()
.execute();
rewrite files 작업에 대한 자세한 내용은 RewriteDataFilesAction을 참고해주세요.
Flink 테이블 유지보수 스트리밍 모드 (Flink Table Maintenance StreamingMode)
개요 (Overview)
Flink 스트리밍 환경에서의 Apache Iceberg 배포에서, 스냅샷 만료, 작은 파일 컴팩션, 고아 파일 정리를 포함한 자동화된 테이블 유지보수 작업 구현은 최적의 쿼리 성능과 스토리지 효율에 중요해요.
전통적으로 이런 유지보수 연산은 Iceberg Spark Actions를 통해서만 사용할 수 있었고, 전용 Spark 클러스터의 배포와 관리가 필요했어요. 테이블 최적화만을 위한 Spark 인프라 의존성은 상당한 아키텍처 복잡성과 운영 오버헤드를 도입해요.
Apache Iceberg의 TableMaintenance API는 Flink 작업이 유지보수 작업을 네이티브로 실행할 수 있게 해줘요. 기존 스트리밍 파이프라인에 임베드하거나 독립 실행형 Flink 작업으로 배포할 수 있어요. 이는 외부 시스템에 대한 의존성을 제거해서 아키텍처를 간소화하고 운영 비용을 줄이며 자동화 기능을 향상시켜요.
지원되는 기능 (Flink) (Supported Features)
ExpireSnapshots
오래된 스냅샷과 그 파일들을 제거해요. 내부적으로 커밋 시 cleanExpiredFiles(true)를 사용하므로, 만료된 메타데이터/파일이 자동으로 정리돼요.
.add(ExpireSnapshots.builder()
.maxSnapshotAge(Duration.ofDays(7))
.retainLast(10)
.deleteBatchSize(1000))
RewriteDataFiles
작은 파일을 컴팩션해서 파일 크기를 최적화해요. 부분 진행 커밋(partial progress commits)과 실행당 최대 재작성 바이트 제한을 지원해요.
.add(RewriteDataFiles.builder()
.targetFileSizeBytes(256 * 1024 * 1024)
.minFileSizeBytes(32 * 1024 * 1024)
.partialProgressEnabled(true)
.partialProgressMaxCommits(5))
DeleteOrphanFiles
아이스버그 테이블의 어떤 메타데이터 파일에서도 참조되지 않아 "고아(orphaned)"로 간주할 수 있는 파일을 제거하는 데 사용돼요. 테이블 위치에서 그런 파일을 검사해요.
.add(DeleteOrphanFiles.builder()
.minAge(Duration.ofDays(3))
.deleteBatchSize(1000))
잠금 관리 (Lock Management)
TriggerLockFactory는 유지보수 작업을 조정하는 데 필수적이에요. 같은 테이블에 대한 동시 유지보수 연산을 방지하는데, 이는 충돌이나 데이터 손상을 일으킬 수 있어요. 같은 작업의 여러 인스턴스가 충돌할 수 있으므로, 이 잠금 메커니즘은 단일 작업에서도 필요해요.
잠금이 필요한 이유 (Why Locks Are Needed)
- 동시 접근: 여러 Flink 작업이 동시에 유지보수를 시도할 수 있어요
- 데이터 일관성: 테이블당 한 번에 하나의 유지보수 연산만 실행되도록 보장해요
- 리소스 관리: 리소스 충돌과 스케줄링 문제를 방지해요
- 중복 작업 방지: 단 하나의 컴팩션 작업만 스케줄돼도 여러 인스턴스가 같은 연산을 시도해서 중복 작업과 낭비되는 리소스가 생길 수 있어요
지원되는 잠금 타입 (Supported Lock Types)
JDBC 잠금 팩토리 (JDBC Lock Factory)
분산 잠금을 관리하기 위해 데이터베이스 테이블을 사용해요:
Map<String, String> jdbcProps = new HashMap<>();
jdbcProps.put("jdbc.user", "flink");
jdbcProps.put("jdbc.password", "flinkpw");
jdbcProps.put("flink-maintenance.lock.jdbc.init-lock-tables", "true"); // Auto-create lock table if it doesn't exist
TriggerLockFactory lockFactory = new JdbcLockFactory(
"jdbc:postgresql://localhost:5432/iceberg", // JDBC URL
"catalog.db.table", // Lock ID (unique identifier)
jdbcProps // JDBC connection properties
);
ZooKeeper 잠금 팩토리 (ZooKeeper Lock Factory)
분산 잠금을 위해 Apache ZooKeeper를 사용해요:
TriggerLockFactory lockFactory = new ZkLockFactory(
"localhost:2181", // ZooKeeper connection string
"catalog.db.table", // Lock ID (unique identifier)
60000, // sessionTimeoutMs
15000, // connectionTimeoutMs
3000, // baseSleepTimeMs
3 // maxRetries
);
Flink 관리 잠금 (Flink-maintained lock)
Flink 자체 내에서 잠금을 유지 관리해요. 외부 시스템을 구성할 필요가 없어요. 유일한 전제 조건은 주어진 테이블에 대해 병렬 테이블 유지보수 작업이 없다는 것이에요.
빠른 시작 (Quick Start)
다음 예시는 Flink 환경에서 아이스버그 테이블의 자동 유지보수 구현을 보여줘요.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
TableLoader tableLoader = TableLoader.fromCatalog(
CatalogLoader.hive("my_catalog", configuration, properties),
TableIdentifier.of("database", "table")
);
Map<String, String> jdbcProps = new HashMap<>();
jdbcProps.put("jdbc.user", "flink");
jdbcProps.put("jdbc.password", "flinkpw");
// JdbcLockFactory Example
TriggerLockFactory lockFactory = new JdbcLockFactory(
"jdbc:postgresql://localhost:5432/iceberg", // JDBC URL
"catalog.db.table", // Lock ID (unique identifier)
jdbcProps // JDBC connection properties
);
// Option 1: With external lock factory (plan to deprecate this Option since 1.12)
TableMaintenance.forTable(env, tableLoader, lockFactory)
// Option 2: With Flink-managed lock (no external lock required)
TableMaintenance.forTable(env, tableLoader)
.uidSuffix("my-maintenance-job")
.rateLimit(Duration.ofMinutes(10))
.lockCheckDelay(Duration.ofSeconds(10))
.add(ExpireSnapshots.builder()
.scheduleOnCommitCount(10)
.maxSnapshotAge(Duration.ofMinutes(10))
.retainLast(5)
.deleteBatchSize(5)
.parallelism(8))
.add(RewriteDataFiles.builder()
.scheduleOnDataFileCount(10)
.targetFileSizeBytes(128 * 1024 * 1024)
.partialProgressEnabled(true)
.partialProgressMaxCommits(10))
.append();
env.execute("Table Maintenance Job");
구성 옵션 (Configuration Options)
TableMaintenance 빌더 (TableMaintenance Builder)
| 메서드 | 설명 | 기본값 |
|---|---|---|
| uidSuffix(String) | 작업의 고유 식별자 접미사 | Random UUID |
| rateLimit(Duration) | 작업 실행 사이의 최소 간격 | 60 seconds |
| lockCheckDelay(Duration) | 잠금 가용성 확인 지연 | 30 seconds |
| parallelism(int) | 유지보수 작업의 기본 병렬도 | System default |
| maxReadBack(int) | 초기화 중에 확인할 최대 스냅샷 수 | 100 |
유지보수 작업 공통 옵션 (Maintenance Task Common Options)
| 메서드 | 설명 | 기본값 | 타입 |
|---|---|---|---|
| scheduleOnCommitCount(int) | N회 커밋 후 트리거 | 자동 스케줄링 없음 | int |
| scheduleOnDataFileCount(int) | N개의 데이터 파일 후 트리거 | 자동 스케줄링 없음 | int |
| scheduleOnDataFileSize(long) | 총 데이터 파일 크기(바이트) 후 트리거 | 자동 스케줄링 없음 | long |
| scheduleOnPosDeleteFileCount(int) | N개의 position delete 파일 후 트리거 | 자동 스케줄링 없음 | int |
| scheduleOnPosDeleteRecordCount(long) | N개의 position delete 레코드 후 트리거 | 자동 스케줄링 없음 | long |
| scheduleOnEqDeleteFileCount(int) | N개의 equality delete 파일 후 트리거 | 자동 스케줄링 없음 | int |
| scheduleOnEqDeleteRecordCount(long) | N개의 equality delete 레코드 후 트리거 | 자동 스케줄링 없음 | long |
| scheduleOnInterval(Duration) | 시간 간격 후 트리거 | 자동 스케줄링 없음 | Duration |
ExpireSnapshots 구성 (ExpireSnapshots Configuration)
| 메서드 | 설명 | 기본값 | 타입 |
|---|---|---|---|
| maxSnapshotAge(Duration) | 보관할 스냅샷의 최대 수명 | 5 days | Duration |
| retainLast(int) | 보관할 최소 스냅샷 수 | 1 | int |
| deleteBatchSize(int) | 각 배치에서 삭제할 파일 수 | 1000 | int |
| planningWorkerPoolSize(int) | 스냅샷 만료 계획을 위한 워커 스레드 수 | Shared worker pool | int |
| cleanExpiredMetadata(boolean) | 스냅샷 만료 시 만료된 메타데이터 파일 제거 | true | boolean |
RewriteDataFiles 구성 (RewriteDataFiles Configuration)
| 메서드 | 설명 | 기본값 | 타입 |
|---|---|---|---|
| targetFileSizeBytes(long) | 재작성 파일의 목표 크기 | Table property 또는 512MB | long |
| minFileSizeBytes(long) | 컴팩션 대상이 되는 파일의 최소 크기 | 목표 파일 크기의 75% | long |
| maxFileSizeBytes(long) | 컴팩션 대상이 되는 파일의 최대 크기 | 목표 파일 크기의 180% | long |
| minInputFiles(int) | 재작성을 트리거할 최소 파일 수 | 5 | int |
| deleteFileThreshold(int) | 재작성을 강제할 데이터 파일당 최소 delete 파일 수 | Integer.MAX_VALUE | int |
| rewriteAll(boolean) | 임계값과 무관하게 모든 데이터 파일 재작성 | false | boolean |
| maxFileGroupSizeBytes(long) | 파일 그룹의 최대 총 크기 | 107374182400 (100GB) | long |
| maxFilesToRewrite(int) | 이 옵션이 지정되지 않으면 모든 적격 파일이 재작성됨 | null | int |
| partialProgressEnabled(boolean) | 부분 진행 커밋 활성화 | false | boolean |
| partialProgressMaxCommits(int) | partialProgressEnabled가 true일 때 부분 진행에 허용되는 최대 커밋 | 10 | int |
| maxRewriteBytes(long) | 실행당 재작성할 최대 바이트 | Long.MAX_VALUE | long |
| filter(Expression) | 재작성할 파일을 선택하는 필터 표현식 | Expressions.alwaysTrue() | Expression |
| maxFileGroupInputFiles(long) | 파일 그룹 내 허용되는 최대 입력 파일 수 | Long.MAX_VALUE | long |
DeleteOrphanFiles 구성 (DeleteOrphanFiles Configuration)
| 메서드 | 설명 | 기본값 | 타입 |
|---|---|---|---|
| location(string) | 제거 대상 후보 파일의 재귀 나열을 시작할 위치 | Table's location | String |
| usePrefixListing(boolean) | true이면 SupportsPrefixOperations 인터페이스를 통해 프리픽스 기반 파일 나열을 사용해요. 이 플래그를 활성화하면 Table FileIO 구현이 SupportsPrefixOperations를 지원해야 해요. (참고: False로 설정하면 재귀 방식으로 파일 정보를 얻어요. 기본 스토리지가 객체 스토리지면 경로를 얻기 위해 API를 반복 호출해요.) | True | boolean |
| prefixMismatchMode(PrefixMismatchMode) | 위치 프리픽스(스킴/authority)가 불일치할 때의 동작: ERROR - 예외 발생. IGNORE - 아무 조치 없음. DELETE - 파일 삭제. | ERROR | PrefixMismatchMode |
| equalSchemes(Map<String, String>) | 동일하게 간주할 파일 시스템 스킴 매핑. 키는 쉼표 구분 스킴 목록, 값은 스킴 | "s3n"=>"s3","s3a"=>"s3" | Map |
| equalAuthorities(Map<String, String>) | 동일하게 간주할 파일 시스템 authority 매핑. 키는 쉼표 구분 authority 목록, 값은 authority | Empty map | Map |
| minAge(Duration) | 이 타임스탬프 이전에 생성된 고아 파일 제거 | 3 days ago | Duration |
| planningWorkerPoolSize(int) | 스냅샷 만료 계획을 위한 워커 스레드 수 | Shared worker pool | int |
완전한 예시 (Complete Example)
public class TableMaintenanceJob {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60000); // Enable checkpointing
// Configure table loader
TableLoader tableLoader = TableLoader.fromCatalog(
CatalogLoader.hive("my_catalog", configuration),
TableIdentifier.of("database", "table")
);
// Set up JDBC lock factory
Map<String, String> jdbcProps = new HashMap<>();
jdbcProps.put("jdbc.user", "flink");
jdbcProps.put("jdbc.password", "flinkpw");
jdbcProps.put("flink-maintenance.lock.jdbc.init-lock-tables", "true");
TriggerLockFactory lockFactory = new JdbcLockFactory(
"jdbc:postgresql://localhost:5432/iceberg",
"catalog.db.table",
jdbcProps
);
// Set up maintenance with comprehensive configuration
TableMaintenance.forTable(env, tableLoader, lockFactory)
.uidSuffix("production-maintenance")
.rateLimit(Duration.ofMinutes(15))
.lockCheckDelay(Duration.ofSeconds(30))
.parallelism(4)
// Daily snapshot cleanup
.add(ExpireSnapshots.builder()
.maxSnapshotAge(Duration.ofDays(7))
.retainLast(10))
// Continuous file optimization
.add(RewriteDataFiles.builder()
.targetFileSizeBytes(256 * 1024 * 1024)
.minFileSizeBytes(32 * 1024 * 1024)
.scheduleOnDataFileCount(20)
.partialProgressEnabled(true)
.partialProgressMaxCommits(5)
.maxRewriteBytes(2L * 1024 * 1024 * 1024)
.parallelism(6))
// Delete orphans files created more than five days ago
.add(DeleteOrphanFiles.builder()
.minAge(Duration.ofDays(5)))
.append();
env.execute("Iceberg Table Maintenance");
}
}
커밋 후 통합이 있는 IcebergSink (IcebergSink with Post-Commit Integration)
Flink용 Apache Iceberg Sink V2는 addPostCommitTopology(...) 메서드를 사용해서 데이터가 테이블에 커밋된 후 유지보수 작업을 자동으로 실행할 수 있게 해줘요.
DataStream API
빌더 (Builder)
IcebergSink.forRowData(dataStream)
.table(table)
.tableLoader(tableLoader)
.rewriteDataFiles(Map.of(
RewriteDataFilesConfig.MAX_BYTES, "1073741824"))
.expireSnapshots(Map.of(
ExpireSnapshotsConfig.RETAIN_LAST, "5",
ExpireSnapshotsConfig.MAX_SNAPSHOT_AGE_SECONDS, "604800"))
.deleteOrphanFiles(Map.of(
DeleteOrphanFilesConfig.MIN_AGE_SECONDS, "259200"))
.append();
구성 (Config)
모든 유지보수 작업은 문자열 속성으로 구성돼요:
Map<String, String> flinkConf = new HashMap<>();
// Enable maintenance tasks
flinkConf.put("flink-maintenance.rewrite.enabled", "true");
flinkConf.put("flink-maintenance.expire-snapshots.enabled", "true");
flinkConf.put("flink-maintenance.delete-orphan-files.enabled", "true");
// Configure rewrite data files
flinkConf.put("flink-maintenance.rewrite.max-bytes", "1073741824");
// Configure expire snapshots
flinkConf.put("flink-maintenance.expire-snapshots.retain-last", "5");
flinkConf.put("flink-maintenance.expire-snapshots.max-snapshot-age-seconds", "604800");
// Configure delete orphan files
flinkConf.put("flink-maintenance.delete-orphan-files.min-age-seconds", "259200");
// Configure JDBC lock settings (deprecated, lock configuration is no longer required for a single Flink job)
flinkConf.put("flink-maintenance.lock.type", "jdbc");
flinkConf.put("flink-maintenance.lock.jdbc.uri", "jdbc:postgresql://localhost:5432/iceberg");
flinkConf.put("flink-maintenance.lock.lock-id", "catalog.db.table");
IcebergSink.forRowData(dataStream)
.table(table)
.tableLoader(tableLoader)
.setAll(flinkConf)
.append();
SQL 예시 (SQL Examples)
쓰기를 실행하기 전에 SQL로 유지보수와 잠금을 활성화하고 구성할 수 있어요.
-- Enable Iceberg V2 Sink and maintenance tasks
SET 'table.exec.iceberg.use.v2.sink' = 'true';
SET 'flink-maintenance.rewrite.enabled' = 'true';
SET 'flink-maintenance.expire-snapshots.enabled' = 'true';
SET 'flink-maintenance.delete-orphan-files.enabled' = 'true';
-- Configure rewrite data files
SET 'flink-maintenance.rewrite.max-bytes' = '1073741824';
-- Configure expire snapshots
SET 'flink-maintenance.expire-snapshots.retain-last' = '5';
-- Configure delete orphan files
SET 'flink-maintenance.delete-orphan-files.min-age-seconds' = '259200';
-- Configure maintenance lock (JDBC)
SET 'flink-maintenance.lock.type' = 'jdbc';
SET 'flink-maintenance.lock.lock-id' = 'catalog.db.table';
SET 'flink-maintenance.lock.jdbc.uri' = 'jdbc:postgresql://localhost:5432/iceberg';
SET 'flink-maintenance.lock.jdbc.init-lock-tables' = 'true';
-- Now run writes; maintenance will be scheduled post-commit
INSERT INTO db.tbl SELECT ...;
또는 테이블 DDL에서 옵션을 지정해요:
CREATE TABLE db.tbl (
...
) WITH (
'connector' = 'iceberg',
'catalog-name' = 'my_catalog',
'catalog-database' = 'db',
'catalog-table' = 'tbl',
'flink-maintenance.rewrite.enabled' = 'true',
'flink-maintenance.expire-snapshots.enabled' = 'true',
'flink-maintenance.delete-orphan-files.enabled' = 'true',
'flink-maintenance.rewrite.max-bytes' = '1073741824',
'flink-maintenance.expire-snapshots.retain-last' = '5',
'flink-maintenance.delete-orphan-files.min-age-seconds' = '259200',
'flink-maintenance.lock.type' = 'jdbc',
'flink-maintenance.lock.lock-id' = 'catalog.db.table',
'flink-maintenance.lock.jdbc.uri' = 'jdbc:postgresql://localhost:5432/iceberg',
'flink-maintenance.lock.jdbc.init-lock-tables' = 'true'
);
IcebergSink 유지보수 구성 (SQL) (IcebergSink Maintenance Configuration)
이 키들은 SQL(SET 또는 테이블 WITH 옵션) 또는 IcebergSink.Builder.set() / setAll()에서 사용돼요.
활성화 플래그 (Enable Flags)
| 키 | 설명 | 기본값 |
|---|---|---|
| flink-maintenance.rewrite.enabled | 컴팩션(데이터 파일 재작성) 활성화 | false |
| flink-maintenance.expire-snapshots.enabled | 스냅샷 만료 활성화 | false |
| flink-maintenance.delete-orphan-files.enabled | 고아 파일 삭제 활성화 | false |
데이터 파일 재작성 구성 (Rewrite Data Files Configuration)
| 키 | 설명 | 기본값 |
|---|---|---|
| flink-maintenance.rewrite.schedule.commit-count | N회 커밋 후 트리거 | 10 |
| flink-maintenance.rewrite.schedule.data-file-count | N개의 데이터 파일 후 트리거 | 1000 |
| flink-maintenance.rewrite.schedule.data-file-size | 총 데이터 파일 크기(바이트) 후 트리거 | 107374182400 (100GB) |
| flink-maintenance.rewrite.schedule.interval-second | 시간 간격(초) 후 트리거 | 600 |
| flink-maintenance.rewrite.max-bytes | 실행당 재작성할 최대 바이트 | Long.MAX_VALUE |
| flink-maintenance.rewrite.partial-progress.enabled | 부분 진행 커밋 활성화 | false |
| flink-maintenance.rewrite.partial-progress.max-commits | 부분 진행의 최대 커밋 | 10 |
스냅샷 만료 구성 (Expire Snapshots Configuration)
| 키 | 설명 | 기본값 |
|---|---|---|
| flink-maintenance.expire-snapshots.schedule.commit-count | N회 커밋 후 트리거 | 10 |
| flink-maintenance.expire-snapshots.schedule.interval-second | 시간 간격(초) 후 트리거 | 3600 (1 hour) |
| flink-maintenance.expire-snapshots.max-snapshot-age-seconds | 보관할 스냅샷의 최대 수명(초) | Not set |
| flink-maintenance.expire-snapshots.retain-last | 보관할 최소 스냅샷 수 | Not set |
| flink-maintenance.expire-snapshots.delete-batch-size | 만료 파일을 삭제하는 배치 크기 | 1000 |
| flink-maintenance.expire-snapshots.clean-expired-metadata | 만료된 메타데이터(파티션 스펙, 스키마) 제거 | true |
| flink-maintenance.expire-snapshots.planning-worker-pool-size | 계획을 위한 워커 풀 크기 | Shared pool |
고아 파일 삭제 구성 (Delete Orphan Files Configuration)
| 키 | 설명 | 기본값 |
|---|---|---|
| flink-maintenance.delete-orphan-files.schedule.interval-second | 시간 간격(초) 후 트리거 | 3600 (1 hour) |
| flink-maintenance.delete-orphan-files.min-age-seconds | 삭제를 고려할 파일의 최소 수명(초) | 259200 (3 days) |
| flink-maintenance.delete-orphan-files.delete-batch-size | 고아 파일 삭제 배치 크기 | 1000 |
| flink-maintenance.delete-orphan-files.location | 재귀 나열을 시작할 위치 | Table location |
| flink-maintenance.delete-orphan-files.use-prefix-listing | 파일 발견에 프리픽스 나열 사용 | true |
| flink-maintenance.delete-orphan-files.planning-worker-pool-size | 계획을 위한 워커 풀 크기 | Shared pool |
| flink-maintenance.delete-orphan-files.equal-schemes | 동등한 스킴 (형식: s3n=s3,s3a=s3) | s3n=s3,s3a=s3 |
| flink-maintenance.delete-orphan-files.equal-authorities | 동등한 authority (형식: auth1=auth2) | Not set |
| flink-maintenance.delete-orphan-files.prefix-mismatch-mode | 프리픽스 불일치 시 동작: ERROR, IGNORE, DELETE | ERROR |
잠금 구성 (SQL) (Lock Configuration)
이 키들은 SQL(SET 또는 테이블 WITH 옵션)에서 사용되며 유지보수가 활성화된 쓰기에 적용돼요.
- JDBC
| 키 | 설명 | 기본값 |
|---|---|---|
| flink-maintenance.lock.type | jdbc로 설정 | |
| flink-maintenance.lock.lock-id | 테이블별 고유 잠금 ID | |
| flink-maintenance.lock.jdbc.uri | JDBC URI | |
| flink-maintenance.lock.jdbc.init-lock-tables | 잠금 테이블 자동 생성 | false |
- ZooKeeper
| 키 | 설명 | 기본값 |
|---|---|---|
| flink-maintenance.lock.type | zookeeper로 설정 | |
| flink-maintenance.lock.lock-id | 테이블별 고유 잠금 ID | |
| flink-maintenance.lock.zookeeper.uri | ZK 연결 URI | |
| flink-maintenance.lock.zookeeper.session-timeout-ms | 세션 타임아웃(ms) | 60000 |
| flink-maintenance.lock.zookeeper.connection-timeout-ms | 연결 타임아웃(ms) | 15000 |
| flink-maintenance.lock.zookeeper.max-retries | 최대 재시도 | 3 |
| flink-maintenance.lock.zookeeper.base-sleep-ms | 재시도 사이 기본 sleep(ms) | 3000 |
| flink-maintenance.lock.zookeeper.max-sleep-ms | 재시도 사이 최대 sleep 시간(ms). 지수 백오프 지연을 상한으로 제한. | 10000 |
| flink-maintenance.lock.zookeeper.retry-policy | ZooKeeper 클라이언트용 재시도 정책 이름. 지원 값: ONE_TIME, N_TIME, BOUNDED_EXPONENTIAL_BACKOFF, UNTIL_ELAPSED, EXPONENTIAL_BACKOFF. | EXPONENTIAL_BACKOFF |
- 코디네이터 잠금 (COORDINATOR LOCK)
| 키 | 설명 | 기본값 |
|---|---|---|
| flink-maintenance.lock.type | ``로 설정하거나 미설정 |
모범 사례 (Best Practices)
리소스 관리 (Resource Management)
- 유지보수 작업을 위해 전용 슬롯 공유 그룹(slot sharing group) 사용
- 클러스터 리소스를 기반으로 적절한 병렬도 설정
- 내결함성을 위해 체크포인트 활성화
스케줄링 전략 (Scheduling Strategy)
- rateLimit으로 너무 잦은 실행을 피해요
- 쓰기량이 많은 테이블에는 scheduleOnCommitCount 사용
- 세밀한 제어가 필요하면 scheduleOnDataFileCount 사용
성능 튜닝 (Performance Tuning)
- 스토리지 성능에 따라 deleteBatchSize 조정
- 큰 재작성 연산에는 partialProgressEnabled 활성화
- 합리적인 maxRewriteBytes 제한 설정
- 적절한 maxFileGroupSizeBytes 설정은 큰 FileGroup을 더 작은 것으로 나눠 병렬 처리 속도를 높일 수 있어요
문제 해결 (Troubleshooting)
파일 삭제 중 OutOfMemoryError (OutOfMemoryError during file deletion)
시나리오: 유지보수 작업이 단일 배치에서 매우 많은 수의 파일을 삭제하려고 할 때 발생할 수 있어요. 특히 보관 이력이 긴 테이블이나 대량 삭제 후에 그렇죠. 원인: 각 파일 삭제는 메타데이터와 객체 스토어 연산을 포함하고, 이는 함께 상당한 메모리를 소비할 수 있어요. 큰 배치는 이 효과를 증폭시켜 JVM 힙을 고갈시킬 수 있어요. 권장사항: 삭제 중 메모리 사용량을 제한하려면 배치 크기를 줄여요.
.deleteBatchSize(500) // Example: 500 files per batch
잠금 충돌 (Lock conflicts)
시나리오: 다중 작업 또는 고가용성 환경에서 두 개 이상의 Flink 작업이 같은 테이블에 대해 동시에 유지보수를 시도할 수 있어요. 원인: 동시 작업이 같은 분산 잠금을 놓고 경쟁해서 재시도와 가능한 지연을 유발해요. 권장사항: 잠금 확인 지연과 rate limit을 늘려서 실패한 시도가 백오프하고 경합을 줄이도록 해요.
.lockCheckDelay(Duration.ofMinutes(1)) // Wait longer before re-checking lock
.rateLimit(Duration.ofMinutes(10)) // Reduce frequency of task execution
느린 재작성 연산 (Slow rewrite operations)
시나리오: 작은 파일이 많은 큰 테이블은 단일 실행에서 테라바이트 단위의 데이터를 재작성해야 할 수 있고, 이는 사용 가능한 리소스를 압도할 수 있어요. 원인: 제한이 없으면 재작성 작업이 한 번에 모든 적격 파일을 처리하려 해서 긴 실행 시간과 가능한 작업 실패로 이어져요. 권장사항: 부분 진행을 활성화해서 재작성된 파일을 더 작은 배치로 커밋하고, 각 실행에서 재작성되는 최대 데이터를 상한으로 제한해요.
.partialProgressEnabled(true) // Commit progress incrementally
.partialProgressMaxCommits(3) // Allow up to 3 commits per run
.maxRewriteBytes(1L * 1024 * 1024 * 1024) // Limit to ~1GB per run