nahidupa opened a new pull request, #18322:
URL: https://github.com/apache/iceberg/pull/18322

   ## Problem
   
   The coordinator's source partition count was fixed at construction. If the 
source assignment expands while that coordinator remains active, reports from 
the old partitions can admit a full commit before the added partition reports. 
The resulting snapshot and completion event can carry a completeness timestamp 
that does not cover the current assignment.
   
   Refreshing only once per cycle is insufficient: an unavailable or unstable 
description can preserve old eligibility, assignments can change before full 
admission, and a late metadata result can overwrite a callback's invalidation.
   
   ## Changes
   
   - Read a stable, nonempty source-group description through the existing 
owned Admin client before each commit cycle.
   - Pin an immutable member-to-topic-partition assignment snapshot for that 
cycle and refresh its expected partition count.
   - Recheck the complete assignment immediately before admitting a full 
commit, including equal-size partition replacements and ownership changes.
   - Invalidate through `CoordinatorThread` from both assignment and revocation 
callbacks, including callbacks that retain the coordinator.
   - Bracket metadata reads with an atomic revision so a successful late 
response cannot restore invalidated eligibility.
   - Keep buffered responses and the existing timeout-partial path. An unknown 
or changed cycle can commit partially without a full completeness watermark, 
but cannot regain full eligibility until a later verified cycle.
   
   No extra Admin client, source-consumer access from the coordinator thread, 
event schema, configuration, public API, or election/fencing protocol is 
introduced.
   
   ## Partial Commits And Limits
   
   When assignment verification fails or the cycle is invalidated, timeout can 
still commit buffered files. Such a partial snapshot omits 
`kafka.connect.valid-through-ts`, and completion events carry a null timestamp. 
A later cycle with freshly verified assignments can complete fully.
   
   This verifies the assignments observed through full-commit admission. It is 
not a distributed generation fence: assignments can change after admission, and 
already-published table snapshots cannot be retracted. Multi-table commits are 
not atomic. Snapshot-expiration recovery and bounded startup replay are outside 
this fix.
   
   This standalone branch intentionally retains main's existing readiness-count 
mechanism. The distinct expected-partition coverage fix belongs to #17925; this 
PR does not duplicate that unmerged change.
   
   ## Validation
   
   Fresh branch from upstream `main` at 
`48330b8dacab6662242d252b39c8444190979bb2`.
   Standalone head: `216658fc1723875772a849d390f36605ca0610ab`.
   
   - **17** new parameterized regression cases cover expansion, 
unknown/unstable/empty metadata after a successful cycle, failed full 
verification, equal-size topology changes, notifications during metadata reads, 
and retained-coordinator added-only/full opens and non-leader closes.
   - Assertions inspect real Iceberg snapshots and added files, exact control 
checkpoints, commit IDs, and full versus null partial timestamps, including 
later verified recovery.
   - The old-partitions-ready/new-partition-missing regression failed on 
unmodified main with a premature full snapshot and passed after the fix.
   - Disposable negative controls failed behavior assertions: frozen 
constructor count **1/1**, omitted callback effect **3/3**, weakened revision 
bracketing **2/2**, omitted pre-full topology comparison **3/3**. Exact source 
restoration was verified and all 17 positives rerun. Compilation failures were 
not counted as behavioral evidence.
   - Standalone full connector gate: **160 tests, 21 suites**, zero failures, 
errors, or skips. Explicit `spotlessCheck`, Checkstyle main/test, class 
uniqueness, and whitespace checks passed.
   
   ```sh
   ./gradlew -DsparkVersions= -DflinkVersions= -DkafkaVersions=3 \
     :iceberg-kafka-connect:iceberg-kafka-connect:spotlessCheck \
     :iceberg-kafka-connect:iceberg-kafka-connect:check \
     -x integrationTest --rerun-tasks --console=plain
   ```
   
   Tests use Corretto 21.0.3, real in-memory Iceberg tables, Kafka mock 
clients, and controlled callback interleavings. They are not live-broker 
expansion/rebalance tests, parallel-thread stress tests, or local JDK 17 
validation. GitHub CI for the published head is separate.
   
   ## Compatibility And Related Work
   
   #17925 remains contributor-approved but unmerged and was left unchanged. In 
an unpublished integration branch, its expected identity set was derived from 
this fix's verified cycle snapshot, preserving identity coverage and 
expected-member overlap. Matching Admin fixtures were added without weakening 
assertions. That combination passed **169 tests, 21 suites** at 
`ae6b264892887b442bd35abc1c8d00dd37642546`.
   
   The unpublished combination with #18006 and #18012 passed **190 tests, 24 
suites** at `25771fca3aa119c40a52423a67f0346d20f1aa6f`, after wiring the 
recovery/lookup metadata fixtures and the consumer factory overload. Those 
integration adaptations will be needed when the independent changes meet; the 
combined branch is not part of this PR.
   
   #17713 is merged and its channel replay guard is preserved. #17450 is closed 
unmerged and is not a dependency. The three bug fixes remain on separate 
branches and PRs.
   
   ---
   **AI Disclosure**
   - Model: [unknown - human to fill in]
   - Platform/Tool: GitHub Copilot.
   - Human Oversight: unreviewed.
   - Prompt Summary: Implement assignment freshness on a fresh independent 
branch, retain timeout partial commits while blocking unverified full 
watermarks, apply approved readiness-review testing lessons, and validate 
standalone plus unpublished combinations without modifying the approved PR.
   


-- 
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