diff --git a/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/IntegrationTestBase.java b/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/IntegrationTestBase.java index 8c011ef537df..af3dc1f736bc 100644 --- a/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/IntegrationTestBase.java +++ b/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/IntegrationTestBase.java @@ -225,7 +225,7 @@ protected void runTest( flush(); Awaitility.await() - .atMost(Duration.ofSeconds(30)) + .atMost(Duration.ofSeconds(60)) .pollInterval(Duration.ofSeconds(1)) .untilAsserted(() -> assertSnapshotAdded(tableIdentifiers)); } diff --git a/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestIntegrationDynamicTable.java b/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestIntegrationDynamicTable.java index 1d3d71a54152..078e7ff254a1 100644 --- a/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestIntegrationDynamicTable.java +++ b/kafka-connect/kafka-connect-runtime/src/integration/java/org/apache/iceberg/connect/TestIntegrationDynamicTable.java @@ -99,7 +99,7 @@ public void testDynamicRouteWithTopicRewritingSMT() { flush(); Awaitility.await() - .atMost(Duration.ofSeconds(30)) + .atMost(Duration.ofSeconds(60)) .pollInterval(Duration.ofSeconds(1)) .untilAsserted(() -> assertSnapshotAdded(List.of(smtTableId))); diff --git a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/Committer.java b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/Committer.java index c0b813f04056..e312f59a6e09 100644 --- a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/Committer.java +++ b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/Committer.java @@ -54,4 +54,6 @@ default void close(Collection closedPartitions) { } void save(Collection sinkRecords); + + default void configure(Catalog catalog, IcebergSinkConfig config, SinkTaskContext context) {} } diff --git a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/IcebergSinkTask.java b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/IcebergSinkTask.java index 471187172e20..0c6b724168bb 100644 --- a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/IcebergSinkTask.java +++ b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/IcebergSinkTask.java @@ -52,6 +52,7 @@ public void start(Map props) { this.config = new IcebergSinkConfig(props); this.catalog = CatalogUtils.loadCatalog(config); this.committer = CommitterFactory.createCommitter(config); + this.committer.configure(catalog, config, context); } @Override 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 01cf165de66b..2d9e4359272a 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 @@ -22,7 +22,6 @@ import java.util.List; import java.util.Map; import java.util.UUID; -import java.util.stream.Collectors; import org.apache.iceberg.connect.IcebergSinkConfig; import org.apache.iceberg.connect.data.Offset; import org.apache.iceberg.connect.events.AvroUtil; @@ -32,11 +31,13 @@ import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.kafka.clients.admin.Admin; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; 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.utils.Utils; import org.apache.kafka.connect.sink.SinkTaskContext; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -53,6 +54,7 @@ abstract class Channel { private final Admin admin; private final Map controlTopicOffsets = Maps.newHashMap(); private final String producerId; + private final String channelId; Channel( String name, @@ -60,13 +62,21 @@ abstract class Channel { IcebergSinkConfig config, KafkaClientFactory clientFactory, SinkTaskContext context) { + this.channelId = config.connectorName() + "-" + config.taskId() + "-" + name; this.controlTopic = config.controlTopic(); this.connectGroupId = config.connectGroupId(); this.context = context; - String transactionalId = config.transactionalPrefix() + name + config.transactionalSuffix(); + String transactionalId = + "worker".equalsIgnoreCase(name) + ? config.transactionalPrefix() + name + config.transactionalSuffix() + : connectGroupId + "-" + config.connectorName() + "-coord"; + this.producer = clientFactory.createProducer(transactionalId); - this.consumer = clientFactory.createConsumer(consumerGroupId); + this.consumer = + "coordinator".equalsIgnoreCase(name) + ? clientFactory.createConsumer(consumerGroupId, "earliest") + : clientFactory.createConsumer(consumerGroupId); this.admin = clientFactory.createAdmin(); this.producerId = UUID.randomUUID().toString(); @@ -85,12 +95,12 @@ protected void send(List events, Map sourceOffset events.stream() .map( event -> { - LOG.info("Sending event of type: {}", event.type().name()); + LOG.info("Channel {} sending event of type: {}", channelId, event.type().name()); byte[] data = AvroUtil.encode(event); // key by producer ID to keep event order return new ProducerRecord<>(controlTopic, producerId, data); }) - .collect(Collectors.toList()); + .toList(); synchronized (producer) { producer.beginTransaction(); @@ -107,7 +117,7 @@ protected void send(List events, Map sourceOffset try { producer.abortTransaction(); } catch (Exception ex) { - LOG.warn("Error aborting producer transaction", ex); + LOG.warn("Channel {} got error while aborting producer transaction", channelId, ex); } throw e; } @@ -119,21 +129,27 @@ protected void send(List events, Map sourceOffset protected void consumeAvailable(Duration pollDuration) { ConsumerRecords records = consumer.poll(pollDuration); while (!records.isEmpty()) { - records.forEach( - record -> { - // the consumer stores the offsets that corresponds to the next record to consume, - // so increment the record offset by one - controlTopicOffsets.put(record.partition(), record.offset() + 1); - - Event event = AvroUtil.decode(record.value()); - - if (event.groupId().equals(connectGroupId)) { - LOG.debug("Received event of type: {}", event.type().name()); - if (receive(new Envelope(event, record.partition(), record.offset()))) { - LOG.info("Handled event of type: {}", event.type().name()); - } - } - }); + for (ConsumerRecord record : records) { + if (record.offset() < controlTopicOffsets.getOrDefault(record.partition(), 0L)) { + LOG.debug( + "Channel {} skipping already-consumed control-topic record at offset {} for partition {}", + channelId, + record.offset(), + record.partition()); + continue; + } + // the consumer stores the offset of the next record to consume, so increment by one + controlTopicOffsets.put(record.partition(), record.offset() + 1); + + Event event = AvroUtil.decode(record.value()); + + if (event.groupId().equals(connectGroupId)) { + LOG.debug("Channel {} received event of type: {}", channelId, event.type().name()); + if (receive(new Envelope(event, record.partition(), record.offset()))) { + LOG.info("Channel {} handled event of type: {}", channelId, event.type().name()); + } + } + } records = consumer.poll(pollDuration); } } @@ -159,9 +175,9 @@ void start() { } void stop() { - LOG.info("Channel stopping"); - producer.close(); - consumer.close(); - admin.close(); + LOG.info("Channel {} stopping", channelId); + Utils.closeQuietly(producer, channelId + "-producer"); + Utils.closeQuietly(consumer, channelId + "-consumer"); + Utils.closeQuietly(admin, channelId + "-admin"); } } 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 7b2d4a25363d..71d190958596 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 @@ -20,16 +20,20 @@ import java.util.Collection; import java.util.Comparator; -import java.util.concurrent.atomic.AtomicBoolean; +import java.util.List; +import java.util.Set; import org.apache.iceberg.catalog.Catalog; import org.apache.iceberg.connect.Committer; import org.apache.iceberg.connect.IcebergSinkConfig; import org.apache.iceberg.connect.data.SinkWriter; import org.apache.iceberg.relocated.com.google.common.annotations.VisibleForTesting; -import org.apache.kafka.clients.admin.Admin; -import org.apache.kafka.clients.admin.ConsumerGroupDescription; -import org.apache.kafka.clients.admin.MemberDescription; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.InvalidProducerEpochException; +import org.apache.kafka.common.errors.ProducerFencedException; +import org.apache.kafka.common.errors.UnknownProducerIdException; import org.apache.kafka.connect.sink.SinkRecord; import org.apache.kafka.connect.sink.SinkTaskContext; import org.slf4j.Logger; @@ -39,92 +43,39 @@ public class CommitterImpl implements Committer { private static final Logger LOG = LoggerFactory.getLogger(CommitterImpl.class); + // Bounded wait for a coordinator thread to fully exit (release its producer/consumer/admin) when + // stopping it, so a newly elected coordinator does not overlap the old one's clients. + private static final long COORDINATOR_STOP_TIMEOUT_MS = 60_000L; + private CoordinatorThread coordinatorThread; private Worker worker; private Catalog catalog; private IcebergSinkConfig config; private SinkTaskContext context; private KafkaClientFactory clientFactory; - private Collection membersWhenWorkerIsCoordinator; - private final AtomicBoolean isInitialized = new AtomicBoolean(false); private String taskId; + private Consumer sourceConsumer; - private void initialize( - Catalog icebergCatalog, - IcebergSinkConfig icebergSinkConfig, - SinkTaskContext sinkTaskContext) { - if (isInitialized.compareAndSet(false, true)) { - this.catalog = icebergCatalog; - this.config = icebergSinkConfig; - this.context = sinkTaskContext; - this.clientFactory = new KafkaClientFactory(config.kafkaProps()); - this.taskId = config.connectorName() + "-" + config.taskId(); - } - } - - static class TopicPartitionComparator implements Comparator { - - @Override - public int compare(TopicPartition o1, TopicPartition o2) { - int result = o1.topic().compareTo(o2.topic()); - if (result == 0) { - result = Integer.compare(o1.partition(), o2.partition()); - } - return result; - } - } - - @VisibleForTesting - boolean hasLeaderPartition(Collection currentAssignedPartitions) { - ConsumerGroupDescription groupDesc; - try (Admin admin = clientFactory.createAdmin()) { - groupDesc = KafkaUtils.consumerGroupDescription(config.connectGroupId(), admin); - } - - Collection members = groupDesc.members(); - if (containsFirstPartition(members, currentAssignedPartitions)) { - membersWhenWorkerIsCoordinator = members; - return true; - } - - return false; - } + private boolean reconcileNeeded = false; @VisibleForTesting - boolean containsFirstPartition( - Collection members, Collection partitions) { - // Determine the first partition across all members to elect the leader - TopicPartition firstTopicPartition = findFirstTopicPartition(members); - - if (firstTopicPartition == null) { - LOG.warn( - "Committer {} found no partitions assigned across all members, cannot determine leader", - taskId); - return false; - } - - boolean containsFirst = partitions.contains(firstTopicPartition); - if (containsFirst) { - LOG.info( - "Committer {} contains the first partition {}, this task is the leader", - taskId, - firstTopicPartition); - } else { - LOG.debug( - "Committer {} does not contain the first partition {}, not the leader", - taskId, - firstTopicPartition); - } - - return containsFirst; + static TopicPartition leaderPartition(Collection subscribedTopics) { + return subscribedTopics.stream() + .min(Comparator.naturalOrder()) + .map(topic -> new TopicPartition(topic, 0)) + .orElse(null); } - @VisibleForTesting - TopicPartition findFirstTopicPartition(Collection members) { - return members.stream() - .flatMap(member -> member.assignment().topicPartitions().stream()) - .min(new TopicPartitionComparator()) - .orElse(null); + @Override + public void configure( + Catalog icebergCatalog, + IcebergSinkConfig icebergSinkConfig, + SinkTaskContext sinkTaskContext) { + this.catalog = icebergCatalog; + this.config = icebergSinkConfig; + this.context = sinkTaskContext; + this.clientFactory = new KafkaClientFactory(config.kafkaProps()); + this.taskId = config.connectorName() + "-" + config.taskId(); } @Override @@ -143,11 +94,7 @@ public void open( IcebergSinkConfig icebergSinkConfig, SinkTaskContext sinkTaskContext, Collection addedPartitions) { - initialize(icebergCatalog, icebergSinkConfig, sinkTaskContext); - if (hasLeaderPartition(addedPartitions)) { - LOG.info("Committer {} received leader partition. Starting Coordinator.", taskId); - startCoordinator(); - } + reconcileNeeded = true; } @Override @@ -162,28 +109,19 @@ public void close(Collection closedPartitions) { // Always try to stop the worker to avoid duplicates. stopWorker(); - // Defensive: close called without prior initialization (should not happen). - if (!isInitialized.get()) { - LOG.warn("Close unexpectedly called on committer {} without partition assignment", taskId); - return; - } - - // Empty partitions → task was stopped explicitly. Stop coordinator if running. + // Empty partitions → the task is being stopped. Tear the coordinator down here since save() + // will not be called again to reconcile it away. if (closedPartitions.isEmpty()) { LOG.info("Committer {} stopped. Closing coordinator.", taskId); stopCoordinator(); return; } - // Normal close: if leader partition is lost, stop coordinator. - if (hasLeaderPartition(closedPartitions)) { - LOG.info("Committer {} lost leader partition. Stopping coordinator.", taskId); - stopCoordinator(); - } + reconcileNeeded = true; // Reset offsets to last committed to avoid data loss. LOG.info("Seeking to last committed offsets for worker {}.", taskId); - KafkaUtils.seekToLastCommittedOffsets(context); + KafkaUtils.seekToLastCommittedOffsets(sourceConsumer()); } @Override @@ -192,13 +130,32 @@ public void save(Collection sinkRecords) { startWorker(); worker.save(sinkRecords); } + if (reconcileNeeded) { + reconcileLeadership(); + reconcileNeeded = false; + } processControlEvents(); } + private void reconcileLeadership() { + Set subscribedTopics = Sets.newTreeSet(sourceConsumer().subscription()); + TopicPartition leader = leaderPartition(subscribedTopics); + if (leader != null && context.assignment().contains(leader)) { + startCoordinator(subscribedTopics); + } else { + stopCoordinator(); + } + } + private void processControlEvents() { if (coordinatorThread != null && coordinatorThread.isTerminated()) { - throw new NotRunningException( - String.format("Coordinator unexpectedly terminated on committer %s", taskId)); + if (isProducerFenced(coordinatorThread.exception())) { + LOG.warn("Committer {} coordinator was fenced by a newer coordinator; clearing it", taskId); + stopCoordinator(); + } else { + throw new NotRunningException( + String.format("Coordinator unexpectedly terminated on committer %s", taskId)); + } } if (worker != null) { worker.process(); @@ -214,11 +171,21 @@ private void startWorker() { } } - private void startCoordinator() { + private void startCoordinator(Set subscribedTopics) { if (null == this.coordinatorThread) { - LOG.info("Task {} elected leader, starting commit coordinator", taskId); + int topicPartitionCount = 0; + for (String topic : subscribedTopics) { + List partitions = sourceConsumer().partitionsFor(topic); + if (partitions != null) { + topicPartitionCount += partitions.size(); + } + } + LOG.info( + "Task {} elected leader, starting commit coordinator (expecting {} partitions)", + taskId, + topicPartitionCount); Coordinator coordinator = - new Coordinator(catalog, config, membersWhenWorkerIsCoordinator, clientFactory, context); + new Coordinator(catalog, config, topicPartitionCount, clientFactory, context); coordinatorThread = new CoordinatorThread(coordinator); coordinatorThread.start(); } @@ -232,9 +199,55 @@ private void stopWorker() { } private void stopCoordinator() { - if (coordinatorThread != null) { - coordinatorThread.terminate(); - coordinatorThread = null; + CoordinatorThread thread = coordinatorThread; + if (thread == null) { + return; + } + coordinatorThread = null; + + try { + if (!thread.isTerminated()) { + thread.terminate(); + } + } catch (RuntimeException e) { + LOG.warn( + "Committer {}: error signalling coordinator termination, continuing shutdown", taskId, e); + } + + try { + thread.join(COORDINATOR_STOP_TIMEOUT_MS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + LOG.warn( + "Committer {}: interrupted while waiting for the coordinator thread to stop", taskId, e); + return; + } + + if (thread.isAlive()) { + LOG.warn( + "Committer {}: coordinator thread did not stop within {} ms; it will keep shutting down " + + "in the background (a newer coordinator fences its producer)", + taskId, + COORDINATOR_STOP_TIMEOUT_MS); + } + } + + private static boolean isProducerFenced(Throwable cause) { + Throwable current = cause; + for (int depth = 0; current != null && depth < 20; depth++, current = current.getCause()) { + if (current instanceof ProducerFencedException + || current instanceof InvalidProducerEpochException + || current instanceof UnknownProducerIdException) { + return true; + } + } + return false; + } + + Consumer sourceConsumer() { + if (null == sourceConsumer) { + sourceConsumer = KafkaUtils.kafkaConsumer(this.context); } + return sourceConsumer; } } 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 1f8b956a5564..01370bac6b24 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 @@ -24,7 +24,6 @@ import java.io.UncheckedIOException; import java.time.Duration; import java.time.OffsetDateTime; -import java.util.Collection; import java.util.List; import java.util.Map; import java.util.Objects; @@ -56,12 +55,12 @@ import org.apache.iceberg.connect.events.TableReference; import org.apache.iceberg.exceptions.CommitFailedException; import org.apache.iceberg.exceptions.NoSuchTableException; +import org.apache.iceberg.relocated.com.google.common.annotations.VisibleForTesting; import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.iceberg.relocated.com.google.common.collect.Streams; import org.apache.iceberg.relocated.com.google.common.util.concurrent.ThreadFactoryBuilder; import org.apache.iceberg.util.SnapshotUtil; import org.apache.iceberg.util.Tasks; -import org.apache.kafka.clients.admin.MemberDescription; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.sink.SinkTaskContext; import org.slf4j.Logger; @@ -90,7 +89,7 @@ class Coordinator extends Channel { Coordinator( Catalog catalog, IcebergSinkConfig config, - Collection members, + int totalPartitionCount, KafkaClientFactory clientFactory, SinkTaskContext context) { // pass consumer group ID to which we commit low watermark offsets @@ -98,8 +97,7 @@ class Coordinator extends Channel { this.catalog = catalog; this.config = config; - this.totalPartitionCount = - members.stream().mapToInt(desc -> desc.assignment().topicPartitions().size()).sum(); + this.totalPartitionCount = totalPartitionCount; this.snapshotOffsetsProp = String.format( "kafka.connect.offsets.%s.%s", config.controlTopic(), config.connectGroupId()); @@ -168,8 +166,8 @@ private void commit(boolean partialCommit) { return; } - if (!(e instanceof CommitFailedException)) { - // CommitStateUnknownException, ValidationException, ForbiddenException, + if (!isRetryable(e)) { + // CommitStateUnknownException, ValidationException, ForbiddenException, ProducerFenced, // NPE, anything else -- not retryable, terminate immediately throw e; } @@ -195,6 +193,14 @@ private void commit(boolean partialCommit) { } } + @VisibleForTesting + static boolean isRetryable(RuntimeException exception) { + return exception instanceof CommitFailedException + || exception instanceof org.apache.kafka.clients.consumer.CommitFailedException + || exception instanceof org.apache.kafka.common.errors.RebalanceInProgressException + || exception instanceof org.apache.kafka.common.errors.RetriableException; + } + private void doCommit(boolean partialCommit) { Map> commitMap = commitState.tableCommitMap(); OffsetDateTime validThroughTs = commitState.validThroughTs(partialCommit); @@ -277,7 +283,7 @@ private void commitToTable( return minOffset == null || envelope.offset() >= minOffset; }) .map(envelope -> (DataWritten) envelope.event().payload()) - .collect(Collectors.toList()); + .toList(); List dataFiles = payloads.stream() @@ -285,7 +291,7 @@ private void commitToTable( .flatMap(payload -> payload.dataFiles().stream()) .filter(dataFile -> dataFile.recordCount() > 0) .filter(distinctByKey(ContentFile::location)) - .collect(Collectors.toList()); + .toList(); List deleteFiles = payloads.stream() @@ -293,7 +299,7 @@ private void commitToTable( .flatMap(payload -> payload.deleteFiles().stream()) .filter(deleteFile -> deleteFile.recordCount() > 0) .filter(distinctByKey(ContentFile::location)) - .collect(Collectors.toList()); + .toList(); if (terminated) { throw new ConnectException( @@ -437,6 +443,7 @@ void terminate() { throw new ConnectException("Timed out waiting for coordinator shutdown"); } } catch (InterruptedException e) { + Thread.currentThread().interrupt(); throw new ConnectException("Interrupted while waiting for coordinator shutdown", e); } } 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 b1a34d0474d9..edca7e515a07 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,7 @@ */ package org.apache.iceberg.connect.channel; +import java.util.concurrent.atomic.AtomicReference; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -27,6 +28,7 @@ class CoordinatorThread extends Thread { private final Coordinator coordinator; private volatile boolean terminated; + private final AtomicReference exception = new AtomicReference<>(); CoordinatorThread(Coordinator coordinator) { super(THREAD_NAME); @@ -39,6 +41,7 @@ public void run() { coordinator.start(); } catch (Exception e) { LOG.error("Coordinator error during start, exiting thread", e); + exception.set(e); this.terminated = true; } @@ -47,6 +50,7 @@ public void run() { coordinator.process(); } catch (Exception e) { LOG.error("Coordinator error during process, exiting thread", e); + exception.set(e); this.terminated = true; } } @@ -66,4 +70,8 @@ void terminate() { this.terminated = true; coordinator.terminate(); } + + Throwable exception() { + return exception.get(); + } } diff --git a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaClientFactory.java b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaClientFactory.java index fd5d27ae34e2..b2ebed6fd74a 100644 --- a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaClientFactory.java +++ b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaClientFactory.java @@ -18,6 +18,7 @@ */ package org.apache.iceberg.connect.channel; +import java.time.Duration; import java.util.Map; import java.util.UUID; import org.apache.iceberg.relocated.com.google.common.collect.Maps; @@ -46,13 +47,22 @@ Producer createProducer(String transactionalId) { producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, transactionalId); KafkaProducer result = new KafkaProducer<>(producerProps, new StringSerializer(), new ByteArraySerializer()); - result.initTransactions(); + try { + result.initTransactions(); + } catch (RuntimeException e) { + result.close(Duration.ZERO); + throw e; + } return result; } Consumer createConsumer(String consumerGroupId) { + return createConsumer(consumerGroupId, "latest"); + } + + Consumer createConsumer(String consumerGroupId, String defaultAutoOffsetReset) { Map consumerProps = Maps.newHashMap(kafkaProps); - consumerProps.putIfAbsent(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); + consumerProps.putIfAbsent(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, defaultAutoOffsetReset); consumerProps.putIfAbsent(ConsumerConfig.CLIENT_ID_CONFIG, UUID.randomUUID().toString()); consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed"); diff --git a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java index bea1c41539c6..587848b5092a 100644 --- a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java +++ b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java @@ -19,12 +19,7 @@ package org.apache.iceberg.connect.channel; import java.util.Map; -import java.util.concurrent.ExecutionException; import org.apache.iceberg.common.DynFields; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; -import org.apache.kafka.clients.admin.Admin; -import org.apache.kafka.clients.admin.ConsumerGroupDescription; -import org.apache.kafka.clients.admin.DescribeConsumerGroupsResult; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerGroupMetadata; import org.apache.kafka.clients.consumer.OffsetAndMetadata; @@ -41,30 +36,17 @@ class KafkaUtils { private static final String CONTEXT_CLASS_NAME = "org.apache.kafka.connect.runtime.WorkerSinkTaskContext"; - static ConsumerGroupDescription consumerGroupDescription(String consumerGroupId, Admin admin) { - try { - DescribeConsumerGroupsResult result = - admin.describeConsumerGroups(ImmutableList.of(consumerGroupId)); - return result.describedGroups().get(consumerGroupId).get(); - - } catch (InterruptedException | ExecutionException e) { - throw new ConnectException( - "Cannot retrieve members for consumer group: " + consumerGroupId, e); - } - } - static ConsumerGroupMetadata consumerGroupMetadata(SinkTaskContext context) { return kafkaConsumer(context).groupMetadata(); } - static void seekToLastCommittedOffsets(SinkTaskContext context) { - Consumer consumer = kafkaConsumer(context); - if (consumer == null) { + static void seekToLastCommittedOffsets(Consumer kafkaConsumer) { + if (kafkaConsumer == null) { return; } Map committedOffsets = - consumer.committed(consumer.assignment()); + kafkaConsumer.committed(kafkaConsumer.assignment()); if (committedOffsets == null || committedOffsets.isEmpty()) { return; } @@ -73,7 +55,7 @@ static void seekToLastCommittedOffsets(SinkTaskContext context) { (topicPartition, offsetAndMetadata) -> { if (offsetAndMetadata != null) { try { - consumer.seek(topicPartition, offsetAndMetadata.offset()); + kafkaConsumer.seek(topicPartition, offsetAndMetadata.offset()); } catch (IllegalStateException e) { LOG.warn( "Rebalance may have occurred, partition {} lost before seeking", @@ -85,7 +67,7 @@ static void seekToLastCommittedOffsets(SinkTaskContext context) { } @SuppressWarnings("unchecked") - private static Consumer kafkaConsumer(SinkTaskContext context) { + static Consumer kafkaConsumer(SinkTaskContext context) { String contextClassName = context.getClass().getName(); try { return ((Consumer) 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 72a362ceacb0..3a3b1711e6e8 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/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 db78b13ae118..70f21b765210 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 @@ -124,6 +124,7 @@ public void before() { clientFactory = mock(KafkaClientFactory.class); when(clientFactory.createProducer(any())).thenReturn(producer); when(clientFactory.createConsumer(any())).thenReturn(consumer); + when(clientFactory.createConsumer(any(), any())).thenReturn(consumer); when(clientFactory.createAdmin()).thenReturn(admin); } 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 f7440dacbe6b..4810022ee9d6 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 @@ -20,104 +20,31 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; -import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.mockStatic; import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; import java.lang.reflect.Field; import java.util.Collections; import java.util.List; -import java.util.Optional; -import org.apache.iceberg.connect.IcebergSinkConfig; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableSet; -import org.apache.kafka.clients.admin.Admin; -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.junit.jupiter.api.Test; -import org.mockito.MockedStatic; public class TestCommitterImpl { @Test - public void testIsLeader() { - MemberAssignment assignment1 = - new MemberAssignment( - ImmutableSet.of(new TopicPartition("topic1", 0), new TopicPartition("topic2", 1))); - MemberDescription member1 = - new MemberDescription(null, Optional.empty(), null, null, assignment1); - - MemberAssignment assignment2 = - new MemberAssignment( - ImmutableSet.of(new TopicPartition("topic2", 0), new TopicPartition("topic1", 1))); - MemberDescription member2 = - new MemberDescription(null, Optional.empty(), null, null, assignment2); - - List members = ImmutableList.of(member1, member2); - - List leaderAssignments = - ImmutableList.of(new TopicPartition("topic2", 1), new TopicPartition("topic1", 0)); - List nonLeaderAssignments = - ImmutableList.of(new TopicPartition("topic2", 0), new TopicPartition("topic1", 1)); + public void testLeaderPartitionIsPartitionZeroOfSmallestTopic() { + // no subscription yet → no leader partition + assertThat(CommitterImpl.leaderPartition(List.of())).isNull(); - CommitterImpl committer = new CommitterImpl(); - assertThat(committer.containsFirstPartition(members, leaderAssignments)).isTrue(); - assertThat(committer.containsFirstPartition(members, nonLeaderAssignments)).isFalse(); - } + // single topic → partition 0 of that topic + assertThat(CommitterImpl.leaderPartition(List.of("only"))) + .isEqualTo(new TopicPartition("only", 0)); - @Test - public void testHasLeaderPartition() throws NoSuchFieldException, IllegalAccessException { - MemberAssignment assignment1 = - new MemberAssignment( - ImmutableSet.of(new TopicPartition("topic1", 0), new TopicPartition("topic2", 1))); - MemberDescription member1 = - new MemberDescription(null, Optional.empty(), null, null, assignment1); - - MemberAssignment assignment2 = - new MemberAssignment( - ImmutableSet.of(new TopicPartition("topic2", 0), new TopicPartition("topic1", 1))); - MemberDescription member2 = - new MemberDescription(null, Optional.empty(), null, null, assignment2); - - List members = ImmutableList.of(member1, member2); - - List leaderAssignments = - ImmutableList.of(new TopicPartition("topic2", 1), new TopicPartition("topic1", 0)); - List nonLeaderAssignments = - ImmutableList.of(new TopicPartition("topic2", 0), new TopicPartition("topic1", 1)); - - CommitterImpl committer = new CommitterImpl(); - Field configField = CommitterImpl.class.getDeclaredField("config"); - Field clientFactoryField = CommitterImpl.class.getDeclaredField("clientFactory"); - configField.setAccessible(true); - clientFactoryField.setAccessible(true); - - IcebergSinkConfig config = mock(IcebergSinkConfig.class); - when(config.connectGroupId()).thenReturn("test-group"); - configField.set(committer, config); - - KafkaClientFactory clientFactory = mock(KafkaClientFactory.class); - Admin admin = mock(Admin.class); - when(clientFactory.createAdmin()).thenReturn(admin); - clientFactoryField.set(committer, clientFactory); - - try (MockedStatic mockKafkaUtils = mockStatic(KafkaUtils.class)) { - ConsumerGroupDescription consumerGroupDescription = mock(ConsumerGroupDescription.class); - mockKafkaUtils - .when(() -> KafkaUtils.consumerGroupDescription(any(), any())) - .thenReturn(consumerGroupDescription); - - when(consumerGroupDescription.members()).thenReturn(members); - - assertThat(committer.hasLeaderPartition(leaderAssignments)).isTrue(); - assertThat(committer.hasLeaderPartition(nonLeaderAssignments)).isFalse(); - } + // multiple topics → partition 0 of the lexicographically-smallest topic + assertThat(CommitterImpl.leaderPartition(List.of("topicB", "topicA", "topicC"))) + .isEqualTo(new TopicPartition("topicA", 0)); } @Test 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 f29c1e5a5132..7fa1f9cf26f4 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 @@ -188,8 +188,7 @@ public void testCommitBoundedRetry() { when(catalog.loadTable(TABLE_IDENTIFIER)).thenReturn(spiedTable); SinkTaskContext context = mock(SinkTaskContext.class); - Coordinator coordinator = - new Coordinator(catalog, config, ImmutableList.of(), clientFactory, context); + Coordinator coordinator = new Coordinator(catalog, config, 0, clientFactory, context); coordinator.start(); initConsumer(); @@ -234,8 +233,7 @@ public void testCommitCounterResetsOnSuccess() { when(catalog.loadTable(TABLE_IDENTIFIER)).thenReturn(spiedTable); SinkTaskContext context = mock(SinkTaskContext.class); - Coordinator coordinator = - new Coordinator(catalog, config, ImmutableList.of(), clientFactory, context); + Coordinator coordinator = new Coordinator(catalog, config, 0, clientFactory, context); coordinator.start(); initConsumer(); @@ -270,8 +268,7 @@ public void testCommitBoundedRetryWithMultipleThreads() { when(catalog.loadTable(TABLE_IDENTIFIER)).thenReturn(spiedTable); SinkTaskContext context = mock(SinkTaskContext.class); - Coordinator coordinator = - new Coordinator(catalog, config, ImmutableList.of(), clientFactory, context); + Coordinator coordinator = new Coordinator(catalog, config, 0, clientFactory, context); coordinator.start(); initConsumer(); @@ -284,6 +281,61 @@ public void testCommitBoundedRetryWithMultipleThreads() { .hasMessageContaining("concurrent update"); } + @Test + public void testIsRetryableClassifiesCommitAndTransientKafkaErrors() { + // Iceberg optimistic-concurrency failure -> retry + assertThat(Coordinator.isRetryable(new CommitFailedException("occ"))).isTrue(); + // Kafka consumer-commit failures from commitConsumerOffsets() -> retry (must NOT be confused + // with the Iceberg CommitFailedException, which was the original bug) + assertThat( + Coordinator.isRetryable( + new org.apache.kafka.clients.consumer.CommitFailedException("evicted"))) + .isTrue(); + assertThat( + Coordinator.isRetryable( + new org.apache.kafka.common.errors.RebalanceInProgressException())) + .isTrue(); + assertThat(Coordinator.isRetryable(new org.apache.kafka.common.errors.TimeoutException())) + .isTrue(); + // Fencing and unknown errors -> fatal (terminate) + assertThat( + Coordinator.isRetryable( + new org.apache.kafka.common.errors.ProducerFencedException("fenced"))) + .isFalse(); + assertThat(Coordinator.isRetryable(new ValidationException("stale"))).isFalse(); + assertThat(Coordinator.isRetryable(new RuntimeException("boom"))).isFalse(); + } + + @Test + public void testConsumeAvailableSkipsAlreadyConsumedOffsets() { + when(config.commitIntervalMs()).thenReturn(Integer.MAX_VALUE); + when(config.commitTimeoutMs()).thenReturn(Integer.MAX_VALUE); + + SinkTaskContext context = mock(SinkTaskContext.class); + Coordinator coordinator = new Coordinator(catalog, config, 0, clientFactory, context); + coordinator.start(); + initConsumer(); + TopicPartition tp = new TopicPartition(CTL_TOPIC_NAME, 0); + + addControlRecord(0); + addControlRecord(1); + coordinator.process(); + assertThat(coordinator.controlTopicOffsets().get(0)).isEqualTo(2L); + + // An eager -coord rebalance can rewind the fetch position and re-deliver already-consumed + // records; consumeAvailable must skip them (offset < tracked), so the tracked offset never + // regresses and the records are not re-processed. + consumer.seek(tp, 0L); + coordinator.process(); + assertThat(coordinator.controlTopicOffsets().get(0)).isEqualTo(2L); + } + + private void addControlRecord(long offset) { + Event event = new Event(config.connectGroupId(), new StartCommit(UUID.randomUUID())); + consumer.addRecord( + new ConsumerRecord<>(CTL_TOPIC_NAME, 0, offset, "key", AvroUtil.encode(event))); + } + private long nextOffset = 1; private void triggerCommitCycle(Coordinator coordinator) { @@ -342,8 +394,7 @@ private UUID coordinatorTest( when(config.commitTimeoutMs()).thenReturn(Integer.MAX_VALUE); SinkTaskContext context = mock(SinkTaskContext.class); - Coordinator coordinator = - new Coordinator(catalog, config, ImmutableList.of(), clientFactory, context); + Coordinator coordinator = new Coordinator(catalog, config, 0, clientFactory, context); coordinator.start(); // init consumer after subscribe() @@ -399,8 +450,7 @@ public void testPartialCommitFailureMetric() { when(catalog.loadTable(TABLE_IDENTIFIER)).thenReturn(spiedTable); SinkTaskContext context = mock(SinkTaskContext.class); - Coordinator coordinator = - new Coordinator(catalog, config, ImmutableList.of(), clientFactory, context); + Coordinator coordinator = new Coordinator(catalog, config, 0, clientFactory, context); coordinator.start(); initConsumer();