apache / apache/iceberg

Flink Sink V2: Add CommitGate plugin interface for deferred commits

Open
#15,770 0 comments 0 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
9.2k
Forks
3.5k
Avg merge
2d 11h
Merged PRs (30d)
132

Description

### Feature Request / Improvement

Adds a `CommitGate` plugin interface to `IcebergSink` that controls whether committables are emitted downstream or buffered in Flink `ListState`. This enables use cases like pausing commits during catalog maintenance operations while keeping the Flink job running.

### Query Engine

Flink

### Motivation

During catalog operations (schema evolution, partition evolution, table migration), it may be necessary to temporarily pause Iceberg commits. Without this capability, the only option is stopping the Flink job, which causes:
- Data loss risk if state is not cleanly checkpointed
- Lag accumulation during the pause window
- Operational toil coordinating job restarts

This adds a `CommitGate` that is checked in `IcebergWriteAggregator.prepareSnapshotPreBarrier()`. When the gate returns `false`, committables are serialized to `ListState` instead of being emitted. When the gate reopens, all buffered committables are flushed in checkpoint order, preserving exactly-once semantics.

### Changes

- New: `CommitGate.java` -- `@FunctionalInterface` with a single method: `boolean isCommitAllowed(long checkpointId)`
- Modified: `IcebergWriteAggregator` -- accepts optional gate, adds `ListState` for buffering, gate check + buffer/flush logic in `prepareSnapshotPreBarrier()`, state initialization in `initializeState()`
- Modified: `IcebergSink.Builder` -- new `commitGate()` method, passed through to aggregator

### Compatibility

- No behavioral change when the gate is not set (null default)
- Buffered committables are checkpointed in `ListState`, so recovery works correctly
- No changes to public API signatures of existing methods
- Fully backward compatible

### Willingness to contribute

- [x] I can contribute this improvement/feature independently
- [x] I would be willing to contribute this improvement/feature with guidance from the Iceberg community
- [ ] I cannot contribute this improvement/feature at this time

Contributor guide

Open the contributing guide

Research direction

Start with IcebergWriteAggregator, especially prepareSnapshotPreBarrier() and initializeState(), then trace how IcebergSink.Builder passes configuration to the aggregator. Read the new CommitGate.java contract and existing committable state handling. Done means gated committables are buffered and checkpointed in order, flushed when allowed, and unchanged when no gate is configured.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
data, stream-processing
Issue type
Feature
Difficulty
4/5
Estimated time
3-5 days
Activity status
Quiet
Clarity
Clearly specified
Newbie friendliness
55/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.