This is an automated email from the ASF dual-hosted git repository.

laskoviymishka pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg.git


The following commit(s) were added to refs/heads/main by this push:
     new dad714b8f0 kafka connect: use a stable coordinator transactional id to 
fence stale coordinators (#18039)
dad714b8f0 is described below

commit dad714b8f0d8ce5e3d2a5f8e2f487dcb71a32100
Author: Thomas Thornton <[email protected]>
AuthorDate: Fri Oct 2 03:56:11 2026 -0700

    kafka connect: use a stable coordinator transactional id to fence stale 
coordinators (#18039)
    
    * kafka connect: use a stable coordinator transactional id to fence stale 
coordinators
    
    Signed-off-by: Thomas Thornton <[email protected]>
    
    * kafka connect: assert exception message in coordinator fencing tests
    
    Signed-off-by: Thomas Thornton <[email protected]>
    
    * kafka connect: skip transaction abort for fenced producers and improve 
fencing tests
    
    Signed-off-by: Thomas Thornton <[email protected]>
    
    ---------
    
    Signed-off-by: Thomas Thornton <[email protected]>
---
 .../org/apache/iceberg/connect/TestContext.java    |  13 ++
 .../iceberg/connect/TestCoordinatorFencing.java    | 164 +++++++++++++++++++++
 .../apache/iceberg/connect/IcebergSinkConfig.java  |  15 +-
 .../apache/iceberg/connect/channel/Channel.java    |  45 +++++-
 .../iceberg/connect/channel/CommitterImpl.java     |  10 +-
 .../iceberg/connect/channel/Coordinator.java       |   7 +-
 .../iceberg/connect/channel/CoordinatorThread.java |  18 +++
 .../connect/channel/NotRunningException.java       |   4 +
 .../org/apache/iceberg/connect/channel/Worker.java |   2 +-
 .../iceberg/connect/channel/ChannelTestBase.java   |  13 ++
 .../iceberg/connect/channel/TestChannel.java       |   3 +-
 .../iceberg/connect/channel/TestCommitterImpl.java |  40 ++++-
 .../iceberg/connect/channel/TestCoordinator.java   | 111 ++++++++------
 13 files changed, 381 insertions(+), 64 deletions(-)

diff --git 
a/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestContext.java
 
b/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestContext.java
index 91828e5c6d..1d3ad9bd93 100644
--- 
a/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestContext.java
+++ 
b/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestContext.java
@@ -120,6 +120,19 @@ public class TestContext {
         new StringSerializer());
   }
 
+  public KafkaProducer<String, String> initLocalTransactionalProducer(String 
transactionalId) {
+    return new KafkaProducer<>(
+        ImmutableMap.of(
+            ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
+            BOOTSTRAP_SERVERS,
+            ProducerConfig.CLIENT_ID_CONFIG,
+            UUID.randomUUID().toString(),
+            ProducerConfig.TRANSACTIONAL_ID_CONFIG,
+            transactionalId),
+        new StringSerializer(),
+        new StringSerializer());
+  }
+
   public Admin initLocalAdmin() {
     return Admin.create(
         ImmutableMap.of(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, 
BOOTSTRAP_SERVERS));
diff --git 
a/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestCoordinatorFencing.java
 
b/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestCoordinatorFencing.java
new file mode 100644
index 0000000000..336b6420e5
--- /dev/null
+++ 
b/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestCoordinatorFencing.java
@@ -0,0 +1,164 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.iceberg.connect;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+import java.time.Duration;
+import java.util.Map;
+import java.util.UUID;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList;
+import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
+import org.apache.kafka.clients.admin.Admin;
+import org.apache.kafka.clients.admin.NewTopic;
+import org.apache.kafka.clients.consumer.ConsumerGroupMetadata;
+import org.apache.kafka.clients.consumer.OffsetAndMetadata;
+import org.apache.kafka.clients.producer.KafkaProducer;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.errors.ProducerFencedException;
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Verifies that a stale coordinator cannot overwrite offsets committed by a 
newer coordinator, even
+ * when the stale coordinator commits an offset ahead of its own previous one 
but still behind the
+ * newer coordinator's.
+ */
+public class TestCoordinatorFencing {
+
+  private static final long STALE_INITIAL_OFFSET = 100L;
+  private static final long NEW_COORDINATOR_OFFSET = 200L;
+  private static final long NEXT_STALE_OFFSET = 150L;
+
+  private final TestContext context = TestContext.instance();
+
+  private String topicName;
+  private String groupId;
+  private TopicPartition topicPartition;
+  private Admin admin;
+
+  @BeforeEach
+  public void before() {
+    topicName = "coord-fencing-topic-" + UUID.randomUUID();
+    groupId = "coord-fencing-group-" + UUID.randomUUID();
+    topicPartition = new TopicPartition(topicName, 0);
+    admin = context.initLocalAdmin();
+    createTopic(topicName);
+  }
+
+  @AfterEach
+  public void after() {
+    deleteTopic(topicName);
+    admin.close();
+  }
+
+  @Test
+  public void newCoordinatorFencesStaleCoordinatorOffsetCommits() throws 
Exception {
+    Map<String, String> connectorProps = connectorProps();
+    String staleCoordinatorId = new 
IcebergSinkConfig(connectorProps).coordinatorTransactionalId();
+    String newCoordinatorId = new 
IcebergSinkConfig(connectorProps).coordinatorTransactionalId();
+    assertThat(newCoordinatorId)
+        .as(
+            
"coordinator·transactional·id·must·be·deterministic·and·match·across"
+                + "independently derived instances for the same connector")
+        .isEqualTo(staleCoordinatorId);
+    KafkaProducer<String, String> staleCoordinator =
+        context.initLocalTransactionalProducer(staleCoordinatorId);
+    KafkaProducer<String, String> newCoordinator =
+        context.initLocalTransactionalProducer(newCoordinatorId);
+
+    try {
+      staleCoordinator.initTransactions();
+      commitOffset(staleCoordinator, STALE_INITIAL_OFFSET);
+      awaitCommittedOffset(STALE_INITIAL_OFFSET);
+
+      newCoordinator.initTransactions();
+      commitOffset(newCoordinator, NEW_COORDINATOR_OFFSET);
+      awaitCommittedOffset(NEW_COORDINATOR_OFFSET);
+
+      assertThatThrownBy(() -> commitOffset(staleCoordinator, 
NEXT_STALE_OFFSET))
+          .isInstanceOf(ProducerFencedException.class)
+          .hasMessageContaining("fence");
+
+      assertThat(committedOffset())
+          .as("fenced coordinator must not clobber the new coordinator's 
committed offset")
+          .isEqualTo(NEW_COORDINATOR_OFFSET);
+    } finally {
+      staleCoordinator.close();
+      newCoordinator.close();
+    }
+  }
+
+  private Map<String, String> connectorProps() {
+    return ImmutableMap.of(
+        "iceberg.catalog.type", "rest",
+        "iceberg.tables", "db.tbl",
+        "name", "coord-fencing-connector-" + UUID.randomUUID());
+  }
+
+  private void commitOffset(KafkaProducer<String, String> producer, long 
offset) {
+    producer.beginTransaction();
+    producer.sendOffsetsToTransaction(
+        ImmutableMap.of(topicPartition, new OffsetAndMetadata(offset)),
+        new ConsumerGroupMetadata(groupId));
+    producer.commitTransaction();
+  }
+
+  private void awaitCommittedOffset(long expected) {
+    Awaitility.await()
+        .atMost(Duration.ofSeconds(30))
+        .pollInterval(Duration.ofMillis(500))
+        .untilAsserted(() -> 
assertThat(committedOffset()).isEqualTo(expected));
+  }
+
+  private long committedOffset() throws InterruptedException, 
ExecutionException, TimeoutException {
+    Map<TopicPartition, OffsetAndMetadata> offsets =
+        admin
+            .listConsumerGroupOffsets(groupId)
+            .partitionsToOffsetAndMetadata()
+            .get(10, TimeUnit.SECONDS);
+    OffsetAndMetadata metadata = offsets.get(topicPartition);
+    return metadata == null ? -1L : metadata.offset();
+  }
+
+  private void createTopic(String topic) {
+    try {
+      admin
+          .createTopics(ImmutableList.of(new NewTopic(topic, 1, (short) 1)))
+          .all()
+          .get(10, TimeUnit.SECONDS);
+    } catch (InterruptedException | ExecutionException | TimeoutException e) {
+      throw new RuntimeException(e);
+    }
+  }
+
+  private void deleteTopic(String topic) {
+    try {
+      admin.deleteTopics(ImmutableList.of(topic)).all().get(10, 
TimeUnit.SECONDS);
+    } catch (InterruptedException | ExecutionException | TimeoutException e) {
+      throw new RuntimeException(e);
+    }
+  }
+}
diff --git 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/IcebergSinkConfig.java
 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/IcebergSinkConfig.java
index fd98b8415f..927d3307f1 100644
--- 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/IcebergSinkConfig.java
+++ 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/IcebergSinkConfig.java
@@ -311,7 +311,7 @@ public class IcebergSinkConfig extends AbstractConfig {
 
   public String transactionalSuffix() {
     // this is for internal use and is not part of the config definition...
-    return originalProps.get(INTERNAL_TRANSACTIONAL_SUFFIX_PROP);
+    return originalProps.getOrDefault(INTERNAL_TRANSACTIONAL_SUFFIX_PROP, "");
   }
 
   public Map<String, String> catalogProps() {
@@ -443,6 +443,19 @@ public class IcebergSinkConfig extends AbstractConfig {
     return "";
   }
 
+  /**
+   * The transactional ID for the coordinator's producer: scoped to the 
coordinator role, stable
+   * across coordinator instances and task restarts within a connector, and 
unique across
+   * connectors. Its stability lets an incoming coordinator's {@code 
initTransactions()} bump the
+   * producer epoch and fence a stale coordinator.
+   *
+   * <p>Unlike the worker producer id, this intentionally omits {@link 
#transactionalSuffix()}: a
+   * per-task/per-instance suffix would make each coordinator's id unique and 
prevent fencing.
+   */
+  public String coordinatorTransactionalId() {
+    return transactionalPrefix() + connectGroupId() + "-coordinator";
+  }
+
   public String hadoopConfDir() {
     return getString(HADOOP_CONF_DIR_PROP);
   }
diff --git 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Channel.java
 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Channel.java
index 0e1f5cb50e..f12febbed9 100644
--- 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Channel.java
+++ 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Channel.java
@@ -38,6 +38,8 @@ import org.apache.kafka.clients.consumer.OffsetAndMetadata;
 import org.apache.kafka.clients.producer.Producer;
 import org.apache.kafka.clients.producer.ProducerRecord;
 import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.errors.InvalidProducerEpochException;
+import org.apache.kafka.common.errors.ProducerFencedException;
 import org.apache.kafka.connect.sink.SinkTaskContext;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -57,8 +59,8 @@ abstract class Channel {
   private final String producerId;
 
   Channel(
-      String name,
       String consumerGroupId,
+      String transactionalId,
       IcebergSinkConfig config,
       KafkaClientFactory clientFactory,
       SinkTaskContext context) {
@@ -66,7 +68,6 @@ abstract class Channel {
     this.connectGroupId = config.connectGroupId();
     this.context = context;
 
-    String transactionalId = config.transactionalPrefix() + name + 
config.transactionalSuffix();
     this.producer = clientFactory.createProducer(transactionalId);
     this.consumer = clientFactory.createConsumer(consumerGroupId);
     this.admin = clientFactory.createAdmin();
@@ -155,10 +156,17 @@ abstract class Channel {
   }
 
   /**
-   * Commit consumer offsets. Only commits offsets if it has not committed 
offsets before or the
-   * value is greater than the cached offset.
+   * Commits consumer offsets in a separate Kafka transaction on the 
coordinator's transactional
+   * producer, committing a partition's offset only when it advances past the 
last committed value.
+   * The producer uses a connector-stable {@code transactional.id}, so a newly 
elected coordinator's
+   * {@code initTransactions()} bumps the producer epoch and fences a 
superseded coordinator, whose
+   * offset commit then fails with a {@link 
org.apache.kafka.common.errors.ProducerFencedException}.
    *
-   * <p>Note: there is a risk that two parallel coordinators may overwrite 
each other's offsets.
+   * <p>This transaction covers only the consumer offset commit, not the 
Iceberg table snapshot
+   * commit. The snapshot commit runs outside any Kafka transaction, so a 
stale coordinator can
+   * still land a snapshot in the window between the new coordinator's {@code 
initTransactions()}
+   * and its own fenced offset commit; that case is guarded separately at the 
Iceberg level by the
+   * {@code SnapshotAncestryValidator} offset validator, not by epoch fencing.
    */
   protected void commitConsumerOffsets() {
     Map<TopicPartition, OffsetAndMetadata> offsetsToCommit = Maps.newHashMap();
@@ -185,13 +193,38 @@ abstract class Channel {
 
     if (!offsetsToCommit.isEmpty()) {
       LOG.debug("Committing consumer offsets: {}", offsetsToCommit);
-      consumer.commitSync(offsetsToCommit);
+      synchronized (producer) {
+        producer.beginTransaction();
+        try {
+          producer.sendOffsetsToTransaction(offsetsToCommit, 
consumer.groupMetadata());
+          producer.commitTransaction();
+        } catch (Exception e) {
+          // fenced producers are fatal and can't abort, so only non-fenced 
producers abort
+          if (!isProducerFenced(e)) {
+            abortTransaction();
+          }
+          throw e;
+        }
+      }
       offsetsToCommit.forEach(
           (topicPartition, metadata) ->
               committedOffsets.put(topicPartition.partition(), 
metadata.offset()));
     }
   }
 
+  private static boolean isProducerFenced(Exception error) {
+    return error instanceof ProducerFencedException
+        || error instanceof InvalidProducerEpochException;
+  }
+
+  private void abortTransaction() {
+    try {
+      producer.abortTransaction();
+    } catch (Exception e) {
+      LOG.warn("Error aborting producer transaction", e);
+    }
+  }
+
   void start() {
     consumer.subscribe(ImmutableList.of(controlTopic));
 
diff --git 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CommitterImpl.java
 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CommitterImpl.java
index 7b2d4a2536..48f044942e 100644
--- 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CommitterImpl.java
+++ 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CommitterImpl.java
@@ -197,8 +197,14 @@ public class CommitterImpl implements Committer {
 
   private void processControlEvents() {
     if (coordinatorThread != null && coordinatorThread.isTerminated()) {
-      throw new NotRunningException(
-          String.format("Coordinator unexpectedly terminated on committer %s", 
taskId));
+      if (coordinatorThread.isFenced()) {
+        LOG.warn("Coordinator on committer {} was fenced by a newer 
coordinator; clearing", taskId);
+        coordinatorThread = null;
+      } else {
+        throw new NotRunningException(
+            String.format("Coordinator unexpectedly terminated on committer 
%s", taskId),
+            coordinatorThread.error());
+      }
     }
     if (worker != null) {
       worker.process();
diff --git 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Coordinator.java
 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Coordinator.java
index 1f8b956a55..0b0778b8fb 100644
--- 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Coordinator.java
+++ 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Coordinator.java
@@ -94,7 +94,12 @@ class Coordinator extends Channel {
       KafkaClientFactory clientFactory,
       SinkTaskContext context) {
     // pass consumer group ID to which we commit low watermark offsets
-    super("coordinator", config.connectGroupId() + "-coord", config, 
clientFactory, context);
+    super(
+        config.connectGroupId() + "-coord",
+        config.coordinatorTransactionalId(),
+        config,
+        clientFactory,
+        context);
 
     this.catalog = catalog;
     this.config = config;
diff --git 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CoordinatorThread.java
 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CoordinatorThread.java
index b1a34d0474..fee7034300 100644
--- 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CoordinatorThread.java
+++ 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CoordinatorThread.java
@@ -18,6 +18,8 @@
  */
 package org.apache.iceberg.connect.channel;
 
+import org.apache.kafka.common.errors.InvalidProducerEpochException;
+import org.apache.kafka.common.errors.ProducerFencedException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -27,6 +29,7 @@ class CoordinatorThread extends Thread {
 
   private final Coordinator coordinator;
   private volatile boolean terminated;
+  private volatile Throwable error;
 
   CoordinatorThread(Coordinator coordinator) {
     super(THREAD_NAME);
@@ -39,6 +42,7 @@ class CoordinatorThread extends Thread {
       coordinator.start();
     } catch (Exception e) {
       LOG.error("Coordinator error during start, exiting thread", e);
+      this.error = e;
       this.terminated = true;
     }
 
@@ -47,6 +51,7 @@ class CoordinatorThread extends Thread {
         coordinator.process();
       } catch (Exception e) {
         LOG.error("Coordinator error during process, exiting thread", e);
+        this.error = e;
         this.terminated = true;
       }
     }
@@ -62,6 +67,19 @@ class CoordinatorThread extends Thread {
     return terminated;
   }
 
+  Throwable error() {
+    return error;
+  }
+
+  /**
+   * Whether the coordinator terminated because a newer coordinator reused its 
{@code
+   * transactional.id} and bumped the producer epoch, fencing this one.
+   */
+  boolean isFenced() {
+    return error instanceof ProducerFencedException
+        || error instanceof InvalidProducerEpochException;
+  }
+
   void terminate() {
     this.terminated = true;
     coordinator.terminate();
diff --git 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/NotRunningException.java
 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/NotRunningException.java
index 72a362ceac..3a3b1711e6 100644
--- 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/NotRunningException.java
+++ 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/NotRunningException.java
@@ -22,4 +22,8 @@ public class NotRunningException extends RuntimeException {
   public NotRunningException(String msg) {
     super(msg);
   }
+
+  public NotRunningException(String msg, Throwable cause) {
+    super(msg, cause);
+  }
 }
diff --git 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Worker.java
 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Worker.java
index 903be70703..722282a4ba 100644
--- 
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Worker.java
+++ 
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/Worker.java
@@ -49,8 +49,8 @@ class Worker extends Channel {
       SinkTaskContext context) {
     // pass transient consumer group ID to which we never commit offsets
     super(
-        "worker",
         config.controlGroupIdPrefix() + UUID.randomUUID(),
+        config.transactionalPrefix() + "worker" + config.transactionalSuffix(),
         config,
         clientFactory,
         context);
diff --git 
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/ChannelTestBase.java
 
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/ChannelTestBase.java
index db78b13ae1..b0767c5066 100644
--- 
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/ChannelTestBase.java
+++ 
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/ChannelTestBase.java
@@ -29,6 +29,7 @@ import static org.mockito.Mockito.when;
 import java.io.IOException;
 import java.util.Collection;
 import java.util.Collections;
+import java.util.Map;
 import org.apache.iceberg.Schema;
 import org.apache.iceberg.Table;
 import org.apache.iceberg.catalog.Namespace;
@@ -39,12 +40,14 @@ import org.apache.iceberg.connect.TableSinkConfig;
 import org.apache.iceberg.inmemory.InMemoryCatalog;
 import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList;
 import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
+import org.apache.iceberg.relocated.com.google.common.collect.Maps;
 import org.apache.iceberg.types.Types;
 import org.apache.kafka.clients.admin.Admin;
 import org.apache.kafka.clients.admin.DescribeTopicsResult;
 import org.apache.kafka.clients.admin.TopicDescription;
 import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
 import org.apache.kafka.clients.consumer.MockConsumer;
+import org.apache.kafka.clients.consumer.OffsetAndMetadata;
 import org.apache.kafka.clients.consumer.OffsetResetStrategy;
 import org.apache.kafka.clients.producer.MockProducer;
 import org.apache.kafka.common.KafkaFuture;
@@ -101,6 +104,8 @@ public class ChannelTestBase {
     when(config.controlTopic()).thenReturn(CTL_TOPIC_NAME);
     when(config.commitThreads()).thenReturn(1);
     when(config.connectGroupId()).thenReturn(CONNECT_CONSUMER_GROUP_ID);
+    when(config.coordinatorTransactionalId())
+        .thenReturn(CONNECT_CONSUMER_GROUP_ID + "-coordinator");
     when(config.tableConfig(any())).thenReturn(mock(TableSinkConfig.class));
     when(config.commitMaxConsecutiveFailures()).thenReturn(1);
 
@@ -138,6 +143,14 @@ public class ChannelTestBase {
     consumer.updateBeginningOffsets(ImmutableMap.of(tp, 0L));
   }
 
+  protected Map<TopicPartition, OffsetAndMetadata> 
committedGroupOffsets(String groupId) {
+    Map<TopicPartition, OffsetAndMetadata> latest = Maps.newHashMap();
+    producer.consumerGroupOffsetsHistory().stream()
+        .filter(commit -> commit.containsKey(groupId))
+        .forEach(commit -> latest.putAll(commit.get(groupId)));
+    return latest;
+  }
+
   private class Listener implements ConsumerRebalanceListener {
     @Override
     public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
diff --git 
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestChannel.java
 
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestChannel.java
index 622a804075..9a198a0f8f 100644
--- 
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestChannel.java
+++ 
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestChannel.java
@@ -29,7 +29,6 @@ import org.apache.iceberg.connect.events.AvroUtil;
 import org.apache.iceberg.connect.events.Event;
 import org.apache.iceberg.connect.events.StartCommit;
 import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
-import org.apache.iceberg.relocated.com.google.common.collect.ImmutableSet;
 import org.apache.iceberg.relocated.com.google.common.collect.Lists;
 import org.apache.kafka.clients.consumer.ConsumerRecord;
 import org.apache.kafka.clients.consumer.OffsetAndMetadata;
@@ -70,7 +69,7 @@ public class TestChannel extends ChannelTestBase {
 
     // the offset committed for the group is what a restarted channel resumes 
from, and what the
     // coordinator stamps on the snapshot, so a regression here is durable
-    assertThat(consumer.committed(ImmutableSet.of(CTL_TOPIC_PARTITION)))
+    assertThat(committedGroupOffsets(consumer.groupMetadata().groupId()))
         .containsEntry(CTL_TOPIC_PARTITION, new OffsetAndMetadata(5L));
   }
 
diff --git 
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestCommitterImpl.java
 
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestCommitterImpl.java
index f7440dacbe..158bf545f1 100644
--- 
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestCommitterImpl.java
+++ 
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestCommitterImpl.java
@@ -40,7 +40,11 @@ import 
org.apache.kafka.clients.admin.ConsumerGroupDescription;
 import org.apache.kafka.clients.admin.MemberAssignment;
 import org.apache.kafka.clients.admin.MemberDescription;
 import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.errors.InvalidProducerEpochException;
+import org.apache.kafka.common.errors.ProducerFencedException;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
 import org.mockito.MockedStatic;
 
 public class TestCommitterImpl {
@@ -123,8 +127,9 @@ public class TestCommitterImpl {
   @Test
   public void testCommitFailurePropagatesAsNotRunningException()
       throws NoSuchFieldException, IllegalAccessException {
+    RuntimeException cause = new RuntimeException("commit failed");
     Coordinator coordinator = mock(Coordinator.class);
-    doThrow(new RuntimeException("commit failed")).when(coordinator).process();
+    doThrow(cause).when(coordinator).process();
 
     CoordinatorThread coordinatorThread = new CoordinatorThread(coordinator);
     coordinatorThread.start();
@@ -132,6 +137,7 @@ public class TestCommitterImpl {
     // wait for the thread to catch the exception, set terminated, and call 
stop
     verify(coordinator, timeout(1000)).stop();
     assertThat(coordinatorThread.isTerminated()).isTrue();
+    assertThat(coordinatorThread.isFenced()).isFalse();
 
     CommitterImpl committer = new CommitterImpl();
     Field field = CommitterImpl.class.getDeclaredField("coordinatorThread");
@@ -140,7 +146,9 @@ public class TestCommitterImpl {
 
     assertThatThrownBy(() -> committer.save(Collections.emptyList()))
         .isInstanceOf(NotRunningException.class)
-        .hasMessageContaining("Coordinator unexpectedly terminated");
+        .hasMessageContaining("Coordinator unexpectedly terminated")
+        .cause()
+        .isSameAs(cause);
   }
 
   @Test
@@ -165,4 +173,32 @@ public class TestCommitterImpl {
         .isInstanceOf(NotRunningException.class)
         .hasMessageContaining("Coordinator unexpectedly terminated");
   }
+
+  @ParameterizedTest
+  @ValueSource(strings = {"ProducerFenced", "InvalidProducerEpoch"})
+  public void testFencedCoordinatorIsClearedWithoutFailingTask(String 
exceptionType)
+      throws NoSuchFieldException, IllegalAccessException {
+    RuntimeException fenceException =
+        "ProducerFenced".equals(exceptionType)
+            ? new ProducerFencedException("fenced by a newer coordinator")
+            : new InvalidProducerEpochException("producer epoch bumped by a 
newer coordinator");
+
+    Coordinator coordinator = mock(Coordinator.class);
+    doThrow(fenceException).when(coordinator).process();
+
+    CoordinatorThread coordinatorThread = new CoordinatorThread(coordinator);
+    coordinatorThread.start();
+
+    verify(coordinator, timeout(1000)).stop();
+    assertThat(coordinatorThread.isTerminated()).isTrue();
+    assertThat(coordinatorThread.isFenced()).isTrue();
+
+    CommitterImpl committer = new CommitterImpl();
+    Field field = CommitterImpl.class.getDeclaredField("coordinatorThread");
+    field.setAccessible(true);
+    field.set(committer, coordinatorThread);
+
+    committer.save(Collections.emptyList());
+    assertThat(field.get(committer)).isNull();
+  }
 }
diff --git 
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestCoordinator.java
 
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestCoordinator.java
index 69e55cca68..9b75d7237e 100644
--- 
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestCoordinator.java
+++ 
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/channel/TestCoordinator.java
@@ -26,9 +26,10 @@ import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.spy;
 import static org.mockito.Mockito.when;
 
+import java.lang.reflect.Field;
 import java.time.OffsetDateTime;
+import java.util.Collections;
 import java.util.List;
-import java.util.Map;
 import java.util.UUID;
 import org.apache.iceberg.AppendFiles;
 import org.apache.iceberg.DataFile;
@@ -55,12 +56,12 @@ import org.apache.iceberg.exceptions.CommitFailedException;
 import org.apache.iceberg.exceptions.ValidationException;
 import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList;
 import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
-import org.apache.iceberg.relocated.com.google.common.collect.ImmutableSet;
 import org.apache.iceberg.relocated.com.google.common.collect.Lists;
 import org.apache.iceberg.types.Types.StructType;
 import org.apache.kafka.clients.consumer.ConsumerRecord;
 import org.apache.kafka.clients.consumer.OffsetAndMetadata;
 import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.errors.ProducerFencedException;
 import org.apache.kafka.connect.sink.SinkTaskContext;
 import org.junit.jupiter.api.Test;
 
@@ -324,8 +325,6 @@ public class TestCoordinator extends ChannelTestBase {
   public void testCommitConsumerOffsetsDoesNotRewind() {
     Coordinator coordinator = startCoordinator();
 
-    TopicPartition ctl = new TopicPartition(CTL_TOPIC_NAME, 0);
-
     long healthyWatermark = 100L;
     coordinator.controlTopicOffsets().put(0, healthyWatermark);
     coordinator.commitConsumerOffsets();
@@ -333,39 +332,62 @@ public class TestCoordinator extends ChannelTestBase {
     coordinator.controlTopicOffsets().put(0, 5L);
     coordinator.commitConsumerOffsets();
 
-    OffsetAndMetadata committedOffsetAndMetadata =
-        consumer.committed(ImmutableSet.of(ctl)).get(ctl);
-    long committed = committedOffsetAndMetadata == null ? 0L : 
committedOffsetAndMetadata.offset();
-
-    assertThat(committed)
-        .as("commitConsumerOffsets should not rewind the shared -coord 
consumer group offsets")
+    assertThat(lastCommittedOffset(0))
+        .as("commitConsumerOffsets should not rewind the -coord consumer group 
offsets")
         .isEqualTo(healthyWatermark);
   }
 
   @Test
-  public void testCommitConsumerDuplicateDoesNotCommit() {
+  public void testFencedCoordinatorCannotCommitOffsets() {
     Coordinator coordinator = startCoordinator();
 
-    TopicPartition ctl = new TopicPartition(CTL_TOPIC_NAME, 0);
-
-    long healthyWatermark = 100L;
-    coordinator.controlTopicOffsets().put(0, healthyWatermark);
+    coordinator.controlTopicOffsets().put(0, 100L);
     coordinator.commitConsumerOffsets();
 
-    long nextWatermark = healthyWatermark + 5;
-    consumer.commitSync(ImmutableMap.of(ctl, new 
OffsetAndMetadata(nextWatermark)));
+    producer.fenceProducer();
 
-    coordinator.controlTopicOffsets().put(0, 100L);
-    coordinator.commitConsumerOffsets();
+    coordinator.controlTopicOffsets().put(0, 150L);
+    assertThatThrownBy(coordinator::commitConsumerOffsets)
+        .as("a fenced coordinator must not be able to commit offsets")
+        .isInstanceOf(ProducerFencedException.class)
+        .hasMessageContaining("fenced");
 
-    OffsetAndMetadata committedOffsetAndMetadata =
-        consumer.committed(ImmutableSet.of(ctl)).get(ctl);
-    long committed = committedOffsetAndMetadata == null ? 0L : 
committedOffsetAndMetadata.offset();
+    assertThat(lastCommittedOffset(0))
+        .as("the fenced coordinator's offset must never reach the offset 
store")
+        .isEqualTo(100L);
+  }
+
+  @Test
+  public void 
testFencedCoordinatorThreadIsClearedByCommitterWithoutFailingTask()
+      throws InterruptedException, NoSuchFieldException, 
IllegalAccessException {
+    when(config.commitIntervalMs()).thenReturn(0);
+    when(config.commitTimeoutMs()).thenReturn(Integer.MAX_VALUE);
 
-    assertThat(committed)
-        .as(
-            "commitConsumerOffsets should not commit offsets when offset has 
not changed relative to local cache")
-        .isEqualTo(nextWatermark);
+    SinkTaskContext context = mock(SinkTaskContext.class);
+    Coordinator coordinator =
+        new Coordinator(catalog, config, ImmutableList.of(), clientFactory, 
context);
+
+    producer.fenceProducer();
+
+    CoordinatorThread coordinatorThread = new CoordinatorThread(coordinator);
+    coordinatorThread.start();
+    coordinatorThread.join(5000);
+
+    assertThat(coordinatorThread.isTerminated()).isTrue();
+    assertThat(coordinatorThread.isFenced())
+        .as("a real coordinator whose producer was fenced must report as 
fenced")
+        .isTrue();
+
+    CommitterImpl committer = new CommitterImpl();
+    Field field = CommitterImpl.class.getDeclaredField("coordinatorThread");
+    field.setAccessible(true);
+    field.set(committer, coordinatorThread);
+
+    committer.save(Collections.emptyList());
+
+    assertThat(field.get(committer))
+        .as("committer must clear a fenced coordinator instead of failing the 
task")
+        .isNull();
   }
 
   @Test
@@ -373,17 +395,10 @@ public class TestCoordinator extends ChannelTestBase {
     Coordinator coordinator = startCoordinator();
 
     long newWatermark = 5L;
-
-    TopicPartition ctl = new TopicPartition(CTL_TOPIC_NAME, 0);
     coordinator.controlTopicOffsets().put(0, newWatermark);
-
     coordinator.commitConsumerOffsets();
 
-    OffsetAndMetadata committedOffsetAndMetadata =
-        consumer.committed(ImmutableSet.of(ctl)).get(ctl);
-    long committed = committedOffsetAndMetadata == null ? 0L : 
committedOffsetAndMetadata.offset();
-
-    assertThat(committed)
+    assertThat(lastCommittedOffset(0))
         .as("commitConsumerOffsets should advance offsets on its first commit")
         .isEqualTo(newWatermark);
   }
@@ -393,8 +408,6 @@ public class TestCoordinator extends ChannelTestBase {
     Coordinator coordinator = startCoordinator();
 
     long healthyWatermark = 100L;
-    TopicPartition ctl = new TopicPartition(CTL_TOPIC_NAME, 0);
-
     coordinator.controlTopicOffsets().put(0, healthyWatermark);
     coordinator.commitConsumerOffsets();
 
@@ -402,11 +415,7 @@ public class TestCoordinator extends ChannelTestBase {
     coordinator.controlTopicOffsets().put(0, watermarkToCommit);
     coordinator.commitConsumerOffsets();
 
-    OffsetAndMetadata committedOffsetAndMetadata =
-        consumer.committed(ImmutableSet.of(ctl)).get(ctl);
-    long committed = committedOffsetAndMetadata == null ? 0L : 
committedOffsetAndMetadata.offset();
-
-    assertThat(committed)
+    assertThat(lastCommittedOffset(0))
         .as("commitConsumerOffsets should advance offsets when its value is 
greater")
         .isEqualTo(watermarkToCommit);
   }
@@ -433,21 +442,25 @@ public class TestCoordinator extends ChannelTestBase {
     coordinator.controlTopicOffsets().put(1, watermarkToSkip);
     coordinator.commitConsumerOffsets();
 
-    Map<TopicPartition, OffsetAndMetadata> committedOffsetAndMetadata =
-        consumer.committed(ImmutableSet.of(ctl0, ctl1));
-
-    OffsetAndMetadata committed0 = committedOffsetAndMetadata.get(ctl0);
-    OffsetAndMetadata committed1 = committedOffsetAndMetadata.get(ctl1);
-
-    assertThat(committed0 == null ? 0L : committed0.offset())
+    assertThat(lastCommittedOffset(0))
         .as("commitConsumerOffsets should advance the consumer group offsets")
         .isEqualTo(watermarkToCommit);
 
-    assertThat(committed1 == null ? 0L : committed1.offset())
+    assertThat(lastCommittedOffset(1))
         .as("commitConsumerOffsets should not rewind consumer group offsets")
         .isEqualTo(healthWatermark1);
   }
 
+  private Long lastCommittedOffset(int partition) {
+    String groupId = consumer.groupMetadata().groupId();
+    OffsetAndMetadata metadata =
+        committedGroupOffsets(groupId).get(new TopicPartition(CTL_TOPIC_NAME, 
partition));
+    assertThat(metadata)
+        .as("expected a committed offset for partition %s in group %s", 
partition, groupId)
+        .isNotNull();
+    return metadata.offset();
+  }
+
   private Coordinator startCoordinator() {
     when(config.commitIntervalMs()).thenReturn(0);
     when(config.commitTimeoutMs()).thenReturn(Integer.MAX_VALUE);

Reply via email to