This is an automated email from the ASF dual-hosted git repository. dsmiley pushed a commit to branch branch_10x in repository https://gitbox.apache.org/repos/asf/solr.git
commit 2d98d6a64c3ff410c557ddacb1b5149e83220f68 Author: Andrzej BiaĆecki <[email protected]> AuthorDate: Wed Sep 2 10:37:03 2026 +0200 SOLR-18408: CrossDC Consumer incorrect nextOffset commits (#4831) (cherry picked from commit de5e1ecd54062e5bbb9491777bbc60bcba20bcdf) --- changelog/unreleased/solr-18408.yml | 8 + .../manager/consumer/KafkaCrossDcConsumer.java | 47 ++-- .../crossdc/manager/consumer/PartitionManager.java | 172 ++++++++++---- .../manager/consumer/KafkaCrossDcConsumerTest.java | 62 ++++- .../manager/consumer/PartitionManagerTest.java | 251 ++++++++++++++++----- .../solr/crossdc/common/KafkaMirroringSink.java | 4 +- .../update/processor/MirroringUpdateProcessor.java | 1 + 7 files changed, 413 insertions(+), 132 deletions(-) diff --git a/changelog/unreleased/solr-18408.yml b/changelog/unreleased/solr-18408.yml new file mode 100644 index 00000000000..85b552c492c --- /dev/null +++ b/changelog/unreleased/solr-18408.yml @@ -0,0 +1,8 @@ +title: SOLR-18408 CrossDC Consumer fix incorrect nextOffset commits + +type: fixed +authors: + - name: Andrzej Bialecki +links: + - name: SOLR-18408 + url: https://issues.apache.org/jira/browse/SOLR-18408 diff --git a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumer.java b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumer.java index 055c16fd9f0..b51122646dc 100644 --- a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumer.java +++ b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumer.java @@ -361,7 +361,7 @@ public class KafkaCrossDcConsumer extends Consumer.CrossDcConsumer { log.trace("Entered pollAndProcessRequests loop"); try { try { - partitionManager.checkOffsetUpdates(); + partitionManager.checkOffsetsAndUpdate(); } catch (Throwable e) { log.error("Error while checking offset updates, shutting down", e); return false; @@ -378,6 +378,9 @@ public class KafkaCrossDcConsumer extends Consumer.CrossDcConsumer { int currentCollapsed = 0; ConsumerRecord<String, MirroredSolrRequest<?>> lastRecord = null; + // the last record merged into updateReqBatch - not the same as lastRecord once a + // subsequent record has triggered a flush of that batch + ConsumerRecord<String, MirroredSolrRequest<?>> batchLastRecord = null; for (TopicPartition partition : records.partitions()) { if (log.isTraceEnabled()) { @@ -388,8 +391,6 @@ public class KafkaCrossDcConsumer extends Consumer.CrossDcConsumer { PartitionManager.PartitionWork partitionWork = partitionManager.getPartitionWork(partition); PartitionManager.WorkUnit workUnit = new PartitionManager.WorkUnit(partition); - workUnit.nextOffset = PartitionManager.getOffsetForPartition(partitionRecords); - partitionWork.partitionQueue.add(workUnit); try { ModifiableSolrParams lastUpdateParams = null; for (ConsumerRecord<String, MirroredSolrRequest<?>> requestRecord : partitionRecords) { @@ -417,19 +418,6 @@ public class KafkaCrossDcConsumer extends Consumer.CrossDcConsumer { if (log.isTraceEnabled()) { log.trace("-- picked type={}, params={}", req.getType(), params); } - if (topicDebug) { - solrReq.addHeader("topic.debug", "true"); - solrReq.addHeader("record.topic", requestRecord.topic()); - solrReq.addHeader("record.partition", String.valueOf(requestRecord.partition())); - solrReq.addHeader("record.offset", String.valueOf(requestRecord.offset())); - solrReq.addHeader("record.timestamp", String.valueOf(requestRecord.timestamp())); - solrReq.addHeader("record.key", requestRecord.key()); - solrReq.addHeader("workUnit.nextOffset", String.valueOf(workUnit.nextOffset)); - solrReq.addHeader("workUnit.partition", String.valueOf(workUnit.partition)); - solrReq.addHeader("workUnit.topic", workUnit.topic); - solrReq.addHeader("workUnit.items", String.valueOf(workUnit.workItems.size())); - } - // determine if it's an UPDATE with deletes, or if the existing batch has deletes boolean hasDeletes = false; if (type == MirroredSolrRequest.Type.UPDATE) { @@ -456,13 +444,27 @@ public class KafkaCrossDcConsumer extends Consumer.CrossDcConsumer { } // send previous batch, if any if (updateReqBatch != null) { - sendBatch(updateReqBatch, type, lastRecord, workUnit); + sendBatch(updateReqBatch, type, batchLastRecord, workUnit); } updateReqBatch = null; currentCollapsed = 0; workUnit = new PartitionManager.WorkUnit(partition); - workUnit.nextOffset = PartitionManager.getOffsetForPartition(partitionRecords); - partitionWork.partitionQueue.add(workUnit); + } + + // this record belongs to the current work unit + partitionWork.assignRecord(workUnit, requestRecord.offset()); + + if (topicDebug) { + solrReq.addHeader("topic.debug", "true"); + solrReq.addHeader("record.topic", requestRecord.topic()); + solrReq.addHeader("record.partition", String.valueOf(requestRecord.partition())); + solrReq.addHeader("record.offset", String.valueOf(requestRecord.offset())); + solrReq.addHeader("record.timestamp", String.valueOf(requestRecord.timestamp())); + solrReq.addHeader("record.key", requestRecord.key()); + solrReq.addHeader("workUnit.nextOffset", String.valueOf(workUnit.nextOffset)); + solrReq.addHeader("workUnit.partition", String.valueOf(workUnit.partition)); + solrReq.addHeader("workUnit.topic", workUnit.topic); + solrReq.addHeader("workUnit.items", String.valueOf(workUnit.workItems.size())); } lastUpdateParams = params; @@ -481,6 +483,7 @@ public class KafkaCrossDcConsumer extends Consumer.CrossDcConsumer { metrics.incrementCollapsedCounter(); currentCollapsed++; } + batchLastRecord = requestRecord; UpdateRequest update = (UpdateRequest) solrReq; MirroredSolrRequest.setParams(updateReqBatch, params); @@ -510,11 +513,11 @@ public class KafkaCrossDcConsumer extends Consumer.CrossDcConsumer { } if (updateReqBatch != null) { - sendBatch(updateReqBatch, MirroredSolrRequest.Type.UPDATE, lastRecord, workUnit); + sendBatch(updateReqBatch, MirroredSolrRequest.Type.UPDATE, batchLastRecord, workUnit); updateReqBatch = null; } try { - partitionManager.checkForOffsetUpdates(partition); + partitionManager.checkOffsetsAndUpdate(partition); } catch (Throwable e) { log.error("Error while checking offset updates, shutting down", e); return false; @@ -542,7 +545,7 @@ public class KafkaCrossDcConsumer extends Consumer.CrossDcConsumer { } try { - partitionManager.checkOffsetUpdates(); + partitionManager.checkOffsetsAndUpdate(); } catch (Throwable e) { log.error("Error while checking offset updates, shutting down", e); return false; diff --git a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java index cef9bfefb98..5504a5e5b78 100644 --- a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java +++ b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java @@ -20,14 +20,13 @@ import com.google.common.annotations.VisibleForTesting; import java.lang.invoke.MethodHandles; import java.util.ArrayDeque; import java.util.HashSet; -import java.util.List; import java.util.Map; import java.util.Queue; import java.util.Set; +import java.util.concurrent.CancellationException; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; -import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; @@ -41,8 +40,52 @@ public class PartitionManager { new ConcurrentHashMap<>(); private final KafkaConsumer<String, MirroredSolrRequest<?>> consumer; - static class PartitionWork { + @VisibleForTesting + public static class PartitionWork { + final TopicPartition partition; final Queue<WorkUnit> partitionQueue = new ArrayDeque<>(); + + PartitionWork(TopicPartition partition) { + this.partition = partition; + } + + /** + * Assign a record to a work unit: enqueue the unit on its first record, and advance its commit + * point to just past that record. A unit that never receives a record is never enqueued, so it + * can never commit an offset of its own. + * + * <p>Guarded by the same monitor as {@link + * PartitionManager#checkOffsetsAndUpdate(TopicPartition)}, which is the other place the queue + * is touched. + * + * @param unit the work unit the record belongs to + * @param recordOffset offset of the record being assigned + * @throws IllegalStateException if recordOffset regresses behind the last record already + * assigned to this unit + */ + synchronized void assignRecord(WorkUnit unit, long recordOffset) { + // does this work unit belong to the partition we're interested in? + if (unit.partition != partition.partition()) { + throw new IllegalStateException( + "Work unit for partition " + + partition.partition() + + " but record for partition " + + unit.partition); + } + if (recordOffset < unit.nextOffset) { + throw new IllegalStateException( + "Out-of-order record offset " + + recordOffset + + ", expected an offset greater than or equal to " + + (unit.nextOffset - 1)); + } + // if this is a new unit enqueue it first + if (unit.nextOffset < 0) { + partitionQueue.add(unit); + } + // advance the commit point to just past the record + unit.nextOffset = recordOffset + 1; + } } @VisibleForTesting @@ -50,7 +93,12 @@ public class PartitionManager { final int partition; final String topic; final Set<Future<?>> workItems = new HashSet<>(); - long nextOffset; + + /** + * Exclusive upper bound of the offsets this unit owns, i.e. the offset to commit once all of + * its work items are done. Negative until the unit is assigned its first record. + */ + long nextOffset = -1; WorkUnit(TopicPartition partition) { this.partition = partition.partition(); @@ -67,56 +115,101 @@ public class PartitionManager { partition, (k, v) -> { if (v == null) { - return new PartitionWork(); + return new PartitionWork(partition); } return v; }); } - public void checkOffsetUpdates() throws Throwable { + public void checkOffsetsAndUpdate() throws Throwable { for (TopicPartition partition : partitionWorkMap.keySet()) { - checkForOffsetUpdates(partition); + checkOffsetsAndUpdate(partition); } } - void checkForOffsetUpdates(TopicPartition partition) throws Throwable { - synchronized (partition) { - PartitionWork work; - if ((work = partitionWorkMap.get(partition)) != null) { - WorkUnit workUnit = work.partitionQueue.peek(); - if (workUnit != null) { - boolean allFuturesDone = true; - for (Future<?> future : workUnit.workItems) { - if (!future.isDone()) { - if (log.isTraceEnabled()) { - log.trace("Future for update is not done topic={}", partition.topic()); - } - allFuturesDone = false; - break; - } + void checkOffsetsAndUpdate(TopicPartition partition) throws Throwable { + // can't synchronize on the argument (equal but distinct object for different threads) + // sync on the PartitionWork instead, which is unique per partition and shared by all threads + // that work on that partition. + final PartitionWork partitionWork = partitionWorkMap.get(partition); + // normally impossible because consumer should always call #getPartitionWork first + // which creates the instance if it doesn't exist. + if (partitionWork == null) { + throw new IllegalStateException( + "PartitionWork for partition " + partition + " not found, likely programming error."); + } - try { - future.get(); - } catch (InterruptedException e) { - log.error("Error updating offset for partition: {}", partition, e); - throw e; - } catch (ExecutionException e) { - log.error("Error updating offset for partition: {}", partition, e); - throw e.getCause(); + synchronized (partitionWork) { + // remove every completed work unit at the head of the queue, stopping at the first one + // that is still in flight - a work unit's offset may only be committed once all of the + // work units before it have been committed too. + long committableOffset = -1; + WorkUnit workUnit; + Throwable failure = null; + try { + while ((workUnit = partitionWork.partitionQueue.peek()) != null) { + if (!isComplete(workUnit, partition)) { + break; + } + // remove completed unit + partitionWork.partitionQueue.poll(); + committableOffset = workUnit.nextOffset; + } + } catch (Throwable t) { + failure = t; + throw t; + } finally { + // commit whatever progress was already verified in this drain, even if a later + // unit's isComplete() threw - otherwise that progress silently gets lost. + if (committableOffset >= 0) { + try { + updateOffset(partition, committableOffset); + } catch (Throwable commitFailure) { + if (commitFailure instanceof InterruptedException) { + Thread.currentThread().interrupt(); } - - if (log.isTraceEnabled()) { - log.trace("Future for update is done topic={}", partition.topic()); + // don't let a secondary commit failure mask the real work-item failure + if (failure != null) { + failure.addSuppressed(commitFailure); + } else { + throw commitFailure; } } + } + } + } + } - if (allFuturesDone) { - work.partitionQueue.poll(); - updateOffset(partition, workUnit.nextOffset); - } + /** Check whether all the work items of this unit are done, rethrowing any of their failures. */ + private boolean isComplete(WorkUnit workUnit, TopicPartition partition) throws Throwable { + for (Future<?> future : workUnit.workItems) { + if (!future.isDone()) { + if (log.isTraceEnabled()) { + log.trace("Future for update is not done topic={}", partition.topic()); } + return false; + } + + try { + // the future is already done, so this returns (or rethrows) without waiting + future.get(); + } catch (InterruptedException e) { + log.error("Error updating offset for partition (interrupted): {}", partition, e); + Thread.currentThread().interrupt(); + throw e; + } catch (CancellationException e) { + log.error("Error updating offset for partition (cancelled): {}", partition, e); + throw e; + } catch (ExecutionException e) { + log.error("Error updating offset for partition: {}", partition, e); + throw e.getCause(); + } + + if (log.isTraceEnabled()) { + log.trace("Future for update is done topic={}", partition.topic()); } } + return true; } /** @@ -136,9 +229,4 @@ public class PartitionManager { consumer.commitSync(Map.of(partition, new OffsetAndMetadata(nextOffset))); } - - static long getOffsetForPartition( - List<ConsumerRecord<String, MirroredSolrRequest<?>>> partitionRecords) { - return partitionRecords.get(partitionRecords.size() - 1).offset() + 1; - } } diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java index 2bc50db48fd..3ce4a2b42d3 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java @@ -229,7 +229,6 @@ public class KafkaCrossDcConsumerTest { @Test public void testHandleFailedResubmit() throws Exception { - // Set up the KafkaCrossDcConsumer KafkaConsumer<String, MirroredSolrRequest<?>> mockConsumer = mock(KafkaConsumer.class); KafkaCrossDcConsumer consumer = createCrossDcConsumerSpy(mockConsumer); @@ -245,9 +244,6 @@ public class KafkaCrossDcConsumerTest { MirroredSolrRequest<?> request = new MirroredSolrRequest<>(new UpdateRequest()); IQueueHandler.Result<MirroredSolrRequest<?>> failedResubmitResult = new IQueueHandler.Result<>(IQueueHandler.ResultStatus.FAILED_RESUBMIT, null, request); - // SolrMessageProcessor mockMessageProcessor = mock(SolrMessageProcessor.class); - // when(mockMessageProcessor.handleItem(any(MirroredSolrRequest.class))) - // .thenReturn(failedResubmitResult); // Mock the KafkaMirroringSink KafkaMirroringSink mockKafkaMirroringSink = mock(KafkaMirroringSink.class); @@ -406,7 +402,7 @@ public class KafkaCrossDcConsumerTest { } // Create a valid MirroredSolrRequest ConsumerRecord<String, MirroredSolrRequest<?>> record = - new ConsumerRecord<>("test-topic", 0, 0, "key", new MirroredSolrRequest<>(validRequest)); + new ConsumerRecord<>("test-topic", 0, i, "key", new MirroredSolrRequest<>(validRequest)); records.add(record); } ConsumerRecords<String, MirroredSolrRequest<?>> consumerRecords = @@ -421,6 +417,62 @@ public class KafkaCrossDcConsumerTest { .sendBatch(any(), eq(MirroredSolrRequest.Type.UPDATE), any(), any()); } + /** + * When a record's differing params force a flush of the batch collapsed so far, the flush must be + * attributed to the last record actually merged into that batch, not to the record that merely + * triggered the flush. + */ + @Test + public void testFlushedBatchLastRecord() { + KafkaConsumer<String, MirroredSolrRequest<?>> mockConsumer = mock(KafkaConsumer.class); + KafkaCrossDcConsumer spyConsumer = createCrossDcConsumerSpy(mockConsumer); + doReturn(new IQueueHandler.Result<>(IQueueHandler.ResultStatus.HANDLED, null)) + .when(messageProcessorMock) + .handleItem(any()); + + UpdateRequest batchRequest1 = new UpdateRequest(); + SolrInputDocument doc1 = new SolrInputDocument(); + doc1.addField("id", "1"); + batchRequest1.add(doc1); + + UpdateRequest batchRequest2 = new UpdateRequest(); + SolrInputDocument doc2 = new SolrInputDocument(); + doc2.addField("id", "2"); + batchRequest2.add(doc2); + + // different params from the first two records, so it can't collapse with them and instead + // forces a flush of the batch they collapsed into + UpdateRequest differentParamsRequest = new UpdateRequest(); + SolrInputDocument doc3 = new SolrInputDocument(); + doc3.addField("id", "3"); + differentParamsRequest.add(doc3); + differentParamsRequest.getParams().set("some.param", "different"); + + ConsumerRecord<String, MirroredSolrRequest<?>> record1 = + new ConsumerRecord<>("test-topic", 0, 0, "key1", new MirroredSolrRequest<>(batchRequest1)); + ConsumerRecord<String, MirroredSolrRequest<?>> record2 = + new ConsumerRecord<>("test-topic", 0, 1, "key2", new MirroredSolrRequest<>(batchRequest2)); + ConsumerRecord<String, MirroredSolrRequest<?>> record3 = + new ConsumerRecord<>( + "test-topic", 0, 2, "key3", new MirroredSolrRequest<>(differentParamsRequest)); + + ConsumerRecords<String, MirroredSolrRequest<?>> records = + new ConsumerRecords<>( + Map.of(new TopicPartition("test-topic", 0), List.of(record1, record2, record3))); + + when(mockConsumer.poll(any())).thenReturn(records).thenThrow(new WakeupException()); + + spyConsumer.run(); + + // record1 and record2 collapsed into one batch; that batch's flush must be attributed to + // record2 (its last record), never to record3 (which only triggered the flush) + verify(spyConsumer, times(1)) + .sendBatch(any(), eq(MirroredSolrRequest.Type.UPDATE), eq(record2), any()); + // record3 starts (and, at end of loop, flushes) its own batch + verify(spyConsumer, times(1)) + .sendBatch(any(), eq(MirroredSolrRequest.Type.UPDATE), eq(record3), any()); + } + @Test public void testHandleInvalidMirroredSolrRequest() { KafkaConsumer<String, MirroredSolrRequest<?>> mockConsumer = mock(KafkaConsumer.class); diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/PartitionManagerTest.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/PartitionManagerTest.java index 91902f5f896..176aa16d276 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/PartitionManagerTest.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/PartitionManagerTest.java @@ -19,13 +19,20 @@ package org.apache.solr.crossdc.manager.consumer; import static org.apache.solr.SolrTestCaseJ4.assumeWorkingMockito; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; +import static org.mockito.ArgumentMatchers.anyMap; +import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; import static org.mockito.Mockito.when; import java.util.Map; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; @@ -35,9 +42,14 @@ import org.apache.kafka.common.TopicPartition; import org.apache.solr.common.util.ExecutorUtil; import org.apache.solr.common.util.SolrNamedThreadFactory; import org.apache.solr.crossdc.common.MirroredSolrRequest; +import org.junit.Before; import org.junit.BeforeClass; import org.junit.Test; +/** + * Tests that a partition's commit point never runs ahead of work that is still in flight: a work + * unit's offset may only be committed once that unit and every unit queued before it are done. + */ @SuppressWarnings("unchecked") public class PartitionManagerTest { @@ -46,15 +58,150 @@ public class PartitionManagerTest { assumeWorkingMockito(); } + private static final TopicPartition PARTITION = new TopicPartition("topic1", 0); + + @SuppressWarnings("unchecked") + private final KafkaConsumer<String, MirroredSolrRequest<?>> consumer = mock(KafkaConsumer.class); + + private PartitionManager partitionManager; + private PartitionManager.PartitionWork work; + + @Before + public void setUp() { + partitionManager = new PartitionManager(consumer); + work = partitionManager.getPartitionWork(PARTITION); + } + + /** Enqueue a work unit owning a single record at the given offset. */ + private PartitionManager.WorkUnit enqueue(long recordOffset) { + PartitionManager.WorkUnit workUnit = new PartitionManager.WorkUnit(PARTITION); + work.assignRecord(workUnit, recordOffset); + return workUnit; + } + + @Test + public void testDrainsAllCompletedUnitsInSingleCommit() throws Throwable { + PartitionManager.WorkUnit first = enqueue(109); + PartitionManager.WorkUnit second = enqueue(119); + PartitionManager.WorkUnit third = enqueue(129); + + // the later units finish first - nothing may be committed while the head is in flight + CompletableFuture<Void> firstWork = new CompletableFuture<>(); + first.workItems.add(firstWork); + second.workItems.add(CompletableFuture.completedFuture(null)); + third.workItems.add(CompletableFuture.completedFuture(null)); + + partitionManager.checkOffsetsAndUpdate(PARTITION); + + verify(consumer, never()).commitSync(anyMap()); + assertEquals(3, work.partitionQueue.size()); + + // once the head completes, all three retire under a single commit of the furthest offset + firstWork.complete(null); + partitionManager.checkOffsetsAndUpdate(PARTITION); + + verify(consumer).commitSync(Map.of(PARTITION, new OffsetAndMetadata(130))); + verifyNoMoreInteractions(consumer); + assertEquals(0, work.partitionQueue.size()); + } + + @Test + public void testStopsAtFirstIncompleteUnit() throws Throwable { + PartitionManager.WorkUnit first = enqueue(109); + PartitionManager.WorkUnit second = enqueue(119); + enqueue(129); + + first.workItems.add(CompletableFuture.completedFuture(null)); + second.workItems.add(new CompletableFuture<>()); + + partitionManager.checkOffsetsAndUpdate(PARTITION); + + // only the first unit's records are done, so only its offset may be committed + verify(consumer).commitSync(Map.of(PARTITION, new OffsetAndMetadata(110))); + verifyNoMoreInteractions(consumer); + assertEquals(2, work.partitionQueue.size()); + assertSame(second, work.partitionQueue.peek()); + } + + @Test + public void testAssignRecordThrowsOnOutOfOrderOffset() { + PartitionManager.WorkUnit unit = new PartitionManager.WorkUnit(PARTITION); + work.assignRecord(unit, 109); + + try { + work.assignRecord(unit, 108); + fail("expected an out-of-order record offset to be rejected"); + } catch (IllegalStateException e) { + // expected + } + } + + @Test + public void testFailedWorkItemPropagatesAndBlocksTheCommit() { + PartitionManager.WorkUnit first = enqueue(109); + first.workItems.add(CompletableFuture.failedFuture(new IllegalStateException("boom"))); + + Throwable thrown = null; + try { + partitionManager.checkOffsetsAndUpdate(PARTITION); + } catch (Throwable e) { + thrown = e; + assertEquals(IllegalStateException.class, e.getClass()); + assertEquals("boom", e.getMessage()); + } + if (thrown == null) { + fail("expected the work item failure to be rethrown"); + } + + verify(consumer, never()).commitSync(anyMap()); + } + + /** + * When a later unit's isComplete() fails after an earlier unit already requires a commit, and the + * commit itself then also fails (e.g. a partition rebalance mid-drain), the original work-item + * failure must still be the one that propagates - with the commit failure attached as suppressed + * rather than replacing the original failure. + */ + @Test + public void testCommitFailurePropgates() { + PartitionManager.WorkUnit first = enqueue(109); + PartitionManager.WorkUnit second = enqueue(119); + + first.workItems.add(CompletableFuture.completedFuture(null)); + second.workItems.add(CompletableFuture.failedFuture(new IllegalStateException("boom"))); + + RuntimeException commitFailure = new RuntimeException("commit failed"); + doThrow(commitFailure).when(consumer).commitSync(anyMap()); + + Throwable thrown = null; + try { + partitionManager.checkOffsetsAndUpdate(PARTITION); + } catch (Throwable e) { + thrown = e; + // the real root cause must still be the one that surfaces + assertEquals(IllegalStateException.class, thrown.getClass()); + assertEquals("boom", thrown.getMessage()); + // with the secondary commit failure preserved rather than discarded + assertEquals(1, thrown.getSuppressed().length); + assertSame(commitFailure, thrown.getSuppressed()[0]); + } + if (thrown == null) { + fail("expected the work item failure to be rethrown"); + } + + // the earlier unit's offset was still attempted, even though the commit itself failed + verify(consumer).commitSync(Map.of(PARTITION, new OffsetAndMetadata(110))); + assertEquals(1, work.partitionQueue.size()); + assertSame(second, work.partitionQueue.peek()); + } + /** * Should return the existing PartitionWork when the partition is already in the partitionWorkMap */ @Test - public void getPartitionWorkWhenPartitionInMap() { - KafkaConsumer<String, MirroredSolrRequest<?>> consumer = mock(KafkaConsumer.class); - PartitionManager partitionManager = new PartitionManager(consumer); + public void testPartitionWorkWhenPartitionInMap() { TopicPartition partition = new TopicPartition("test-topic", 0); - PartitionManager.PartitionWork partitionWork = new PartitionManager.PartitionWork(); + PartitionManager.PartitionWork partitionWork = new PartitionManager.PartitionWork(partition); partitionManager.partitionWorkMap.put(partition, partitionWork); PartitionManager.PartitionWork result = partitionManager.getPartitionWork(partition); @@ -65,9 +212,7 @@ public class PartitionManagerTest { /** Should create a new PartitionWork when the partition is not in the partitionWorkMap */ @Test - public void getPartitionWorkWhenPartitionNotInMap() { - KafkaConsumer<String, MirroredSolrRequest<?>> consumer = mock(KafkaConsumer.class); - PartitionManager partitionManager = new PartitionManager(consumer); + public void testPartitionWorkWhenPartitionNotInMap() { TopicPartition partition = new TopicPartition("test-topic", 0); PartitionManager.PartitionWork partitionWork = partitionManager.getPartitionWork(partition); @@ -79,83 +224,65 @@ public class PartitionManagerTest { /** Should not update the offset when the future for update is not done */ @Test - public void checkForOffsetUpdatesWhenFutureNotDone() throws Throwable { - KafkaConsumer<String, MirroredSolrRequest<?>> consumer = mock(KafkaConsumer.class); - PartitionManager partitionManager = new PartitionManager(consumer); - TopicPartition partition = new TopicPartition("test-topic", 0); - PartitionManager.PartitionWork partitionWork = partitionManager.getPartitionWork(partition); - PartitionManager.WorkUnit workUnit = new PartitionManager.WorkUnit(partition); + public void testForOffsetUpdatesWhenFutureNotDone() throws Throwable { + PartitionManager.WorkUnit workUnit = new PartitionManager.WorkUnit(PARTITION); Future<?> future = mock(Future.class); when(future.isDone()).thenReturn(false); workUnit.workItems.add(future); - partitionWork.partitionQueue.add(workUnit); + work.assignRecord(workUnit, 0); - partitionManager.checkForOffsetUpdates(partition); + partitionManager.checkOffsetsAndUpdate(PARTITION); - assertEquals(1, partitionWork.partitionQueue.size()); - assertTrue(partitionWork.partitionQueue.contains(workUnit)); + assertEquals(1, work.partitionQueue.size()); + assertTrue(work.partitionQueue.contains(workUnit)); } /** Should update the offset when the future for update is done */ @Test - public void checkForOffsetUpdatesWhenFutureDone() throws Throwable { - KafkaConsumer<String, MirroredSolrRequest<?>> consumer = mock(KafkaConsumer.class); - PartitionManager partitionManager = new PartitionManager(consumer); - TopicPartition partition = new TopicPartition("test-topic", 0); - - PartitionManager.PartitionWork partitionWork = partitionManager.getPartitionWork(partition); - PartitionManager.WorkUnit workUnit = new PartitionManager.WorkUnit(partition); - partitionWork.partitionQueue.add(workUnit); + public void testForOffsetUpdatesWhenFutureDone() throws Throwable { + PartitionManager.WorkUnit workUnit = new PartitionManager.WorkUnit(PARTITION); + work.assignRecord(workUnit, 0); // Use a real Future instead of a mocked one ExecutorService executor = ExecutorUtil.newMDCAwareSingleThreadExecutor(new SolrNamedThreadFactory("test")); - Future<?> future = - executor.submit( - () -> { - // Simulate the task being completed - }); + try { + Future<?> future = + executor.submit( + () -> { + // Simulate the task being completed + }); - workUnit.workItems.add(future); - - // Wait for the Future to completeE - future.get(10, TimeUnit.SECONDS); + workUnit.workItems.add(future); - partitionManager.checkForOffsetUpdates(partition); + // Wait for the Future to completeE + future.get(10, TimeUnit.SECONDS); - // Verify that the consumer.commitSync() method was called with the correct parameters - verify(consumer, times(1)) - .commitSync(Map.of(partition, new OffsetAndMetadata(workUnit.nextOffset))); + partitionManager.checkOffsetsAndUpdate(PARTITION); - // Verify that the partitionQueue is empty after processing - assertTrue(partitionWork.partitionQueue.isEmpty()); + // Verify that the consumer.commitSync() method was called with the correct parameters + verify(consumer, times(1)) + .commitSync(Map.of(PARTITION, new OffsetAndMetadata(workUnit.nextOffset))); - // Shutdown the executor - executor.shutdown(); + // Verify that the partitionQueue is empty after processing + assertTrue(work.partitionQueue.isEmpty()); + } finally { + executor.shutdown(); + } } /** Should check for offset updates for all partitions in the partitionWorkMap */ @Test - public void checkOffsetUpdatesForAllPartitions() throws Throwable { // Create a mock KafkaConsumer - KafkaConsumer<String, MirroredSolrRequest<?>> mockConsumer = mock(KafkaConsumer.class); - - // Create a PartitionManager instance with the mock KafkaConsumer - PartitionManager partitionManager = new PartitionManager(mockConsumer); - - // Create a few TopicPartitions - TopicPartition partition1 = new TopicPartition("topic1", 0); + public void testOffsetUpdatesForAllPartitions() throws Throwable { + // reuse the shared partition for the first partition, and add a second one TopicPartition partition2 = new TopicPartition("topic2", 0); - - // Add some PartitionWork to the partitionWorkMap - PartitionManager.PartitionWork work1 = partitionManager.getPartitionWork(partition1); PartitionManager.PartitionWork work2 = partitionManager.getPartitionWork(partition2); - // Create WorkUnits and add them to the PartitionWork - PartitionManager.WorkUnit workUnit1 = new PartitionManager.WorkUnit(partition1); + PartitionManager.WorkUnit workUnit1 = new PartitionManager.WorkUnit(PARTITION); PartitionManager.WorkUnit workUnit2 = new PartitionManager.WorkUnit(partition2); - work1.partitionQueue.add(workUnit1); - work2.partitionQueue.add(workUnit2); + work.assignRecord(workUnit1, 0); + work2.assignRecord(workUnit2, 0); // Create mock Futures and add them to the WorkUnits Future<?> mockFuture1 = mock(Future.class); @@ -168,17 +295,17 @@ public class PartitionManagerTest { when(mockFuture1.isDone()).thenReturn(true); when(mockFuture2.isDone()).thenReturn(true); - // Call the checkOffsetUpdates method - partitionManager.checkOffsetUpdates(); + // Call the checkOffsetsAndUpdate method + partitionManager.checkOffsetsAndUpdate(); // Verify that the futures were checked for completion verify(mockFuture1, times(1)).isDone(); verify(mockFuture2, times(1)).isDone(); // Verify that the updateOffset method was called for each partition - verify(mockConsumer, times(1)) - .commitSync(Map.of(partition1, new OffsetAndMetadata(workUnit1.nextOffset))); - verify(mockConsumer, times(1)) + verify(consumer, times(1)) + .commitSync(Map.of(PARTITION, new OffsetAndMetadata(workUnit1.nextOffset))); + verify(consumer, times(1)) .commitSync(Map.of(partition2, new OffsetAndMetadata(workUnit2.nextOffset))); } } diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java index c348f032c21..8935160937f 100644 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java +++ b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java @@ -34,6 +34,7 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; +import org.apache.solr.common.util.IOUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -118,7 +119,7 @@ public class KafkaMirroringSink implements RequestMirroringSink, Closeable { slowSubmitAction(elapsedTimeMillis); } } catch (Exception e) { - // We are intentionally catching all exceptions, the expected exception form this function is + // We are intentionally catching all exceptions, the expected exception from this function is // {@link MirroringException} String message = "Unable to enqueue request " @@ -232,5 +233,6 @@ public class KafkaMirroringSink implements RequestMirroringSink, Closeable { producer.flush(); producer.close(); } + IOUtils.closeQuietly(consumer); } } diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java index 284dd91ca7d..223616806d6 100644 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java +++ b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java @@ -273,6 +273,7 @@ public class MirroringUpdateProcessor extends UpdateRequestProcessor { } catch (Exception e) { log.error("mirror submit failed", e); producerMetrics.getSubmittedDeleteByIdError().inc(); + producerMetrics.getSubmitError().inc(); throw new SolrException(SERVER_ERROR, "mirror submit failed", e); } }
