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]

Reply via email to