kumarpritam863 opened a new pull request, #17450:
URL: https://github.com/apache/iceberg/pull/17450
## What
Reworks how the Iceberg Kafka Connect sink elects its single commit
`Coordinator` and
hardens the coordinator's fencing, recovery, and shutdown paths. No
control-topic
wire-format change; exactly-once semantics are preserved.
## Leader election — level read, no Admin call
- **Before:** the coordinator was chosen on every rebalance by calling
`Admin.describeConsumerGroups(connectGroupId)` and picking the task that
owned the
globally-lowest `(topic, partition)` from a transient, mid-rebalance
member snapshot.
This required a `DESCRIBE` ACL, added an Admin round-trip to the
rebalance path, and was
racy during cooperative rebalancing.
- **After:** the leader is the task whose `assignment()` contains
**partition 0 of the
lexicographically-smallest topic in its own `subscription()`** —
connector-wide,
identical across tasks, and valid for both a `topics` list and a
`topics.regex`. No Admin
call.
- **Leadership is read as a level on the task thread.** `open`/`close`
only flag that a
reconcile is needed; `save()` reconciles once — start the coordinator if
this task owns
the leader partition, stop it otherwise (both idempotent). Keeping
start/stop off the
rebalance callback avoids blocking RPCs there and stops an **eager**
rebalance from
needlessly restarting a still-leading coordinator (which would discard
in-flight commit
state).
- The coordinator's commit-readiness partition count now comes from
`consumer.partitionsFor(...)` over the subscribed topics instead of a
member snapshot.
- Adds a `Committer.configure(Catalog, IcebergSinkConfig,
SinkTaskContext)` lifecycle hook
(invoked from `IcebergSinkTask.start`) for one-time setup.
## Coordinator hardening
- **Zombie fencing via a fixed `transactional.id`.** The coordinator
producer id is now
`<connectGroupId>-<connectorName>-coord` — identical across a
connector's tasks and
stable across restarts — so a newly elected coordinator's
`initTransactions()`
epoch-fences a prior (zombie) coordinator's control-plane writes. Worker
ids are
unchanged. This fences the brief two-coordinator overlap the level-read
handoff can
create (the losing task stops on its *next* `save()`).
- **Fenced ≠ fatal.** A coordinator that terminates because it was fenced
(`ProducerFenced` / `InvalidProducerEpoch` / `UnknownProducerId`,
matched across the
cause chain) is cleared without failing the task; any other termination
still fails the
task.
- **Recovery reads from `earliest`.** The `-coord` consumer group defaults
to
`auto.offset.reset=earliest`, so a fresh or expired group re-reads
uncommitted control
events instead of skipping to the log end. Replay is idempotent via the
snapshot offset
floor + `distinctByKey(location)` dedup + the offsets compare-and-swap.
- **Idempotency filter in `Channel.consumeAvailable`.** Control-topic
records at or below
the already-consumed per-partition offset are skipped, so a re-delivered
or rewound
record is never re-buffered or re-counted.
- **Transient Kafka commit errors are retryable.** Errors from
`commitConsumerOffsets`
(Iceberg `CommitFailedException`, Kafka `CommitFailedException`,
`RebalanceInProgressException`, `RetriableException`) are retried rather
than failing the
task — the table commit already succeeded.
- **Bounded, interrupt-safe shutdown.** `stopCoordinator` clears state
first, then
bounded-joins the thread; `terminate()` failures are best-effort and
interrupts are
preserved.
- **No producer leak on failed init.** `KafkaClientFactory.createProducer`
closes the
producer if `initTransactions()` throws.
## Compatibility
- No config or control-topic wire-format changes.
- Exactly-once unchanged — anchored on the Iceberg offsets
compare-and-swap +
file-location dedup; the fixed `transactional.id` adds control-plane
fencing on top.
## Testing
- Unit tests for the election key (`leaderPartition`), the
retryable-commit classification,
and coordinator fenced-vs-fatal termination.
- Integration suite passes (13 tests).
- Raised the integration-test commit-wait from 30s to 60s to reduce
CI-load flakiness.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]