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