This is an automated email from the ASF dual-hosted git repository.
dsmiley pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/solr.git
The following commit(s) were added to refs/heads/main by this push:
new 0cc72b8e326 SOLR-18406: Restart replication after index generation
expires (#4820)
0cc72b8e326 is described below
commit 0cc72b8e326d9702958e9603b420af4c60d24c19
Author: Zhenyu Li <[email protected]>
AuthorDate: Mon Sep 7 23:04:40 2026 -0400
SOLR-18406: Restart replication after index generation expires (#4820)
Leader/follower replication now restarts with the latest index generation
when the selected generation expires during download, avoiding retries against
an unavailable generation.
---
.../SOLR-18406-replication-stale-generation.yml | 10 ++++
.../java/org/apache/solr/handler/IndexFetcher.java | 65 +++++++++++++++++-----
.../solr/handler/admin/api/ReplicationAPIBase.java | 11 +++-
.../solr/handler/TestReplicationHandler.java | 62 +++++++++++++++++++++
4 files changed, 133 insertions(+), 15 deletions(-)
diff --git a/changelog/unreleased/SOLR-18406-replication-stale-generation.yml
b/changelog/unreleased/SOLR-18406-replication-stale-generation.yml
new file mode 100644
index 00000000000..694dd5c4b08
--- /dev/null
+++ b/changelog/unreleased/SOLR-18406-replication-stale-generation.yml
@@ -0,0 +1,10 @@
+title: >
+ Leader/follower replication now restarts with the latest index generation
when the selected
+ generation expires during download, avoiding retries against an unavailable
generation.
+type: changed
+authors:
+ - name: ZhenyuLi
+ nick: JHSUYU
+links:
+ - name: SOLR-18406
+ url: https://issues.apache.org/jira/browse/SOLR-18406
diff --git a/solr/core/src/java/org/apache/solr/handler/IndexFetcher.java
b/solr/core/src/java/org/apache/solr/handler/IndexFetcher.java
index 4a9f14834cd..b6fa1146959 100644
--- a/solr/core/src/java/org/apache/solr/handler/IndexFetcher.java
+++ b/solr/core/src/java/org/apache/solr/handler/IndexFetcher.java
@@ -405,6 +405,18 @@ public class IndexFetcher {
*/
IndexFetchResult fetchLatestIndex(boolean forceReplication, boolean
forceCoreReload)
throws IOException, InterruptedException {
+ try {
+ return fetchLatestIndexOnce(forceReplication, forceCoreReload);
+ } catch (InvalidIndexGenerationException e) {
+ log.info(
+ "Leader no longer has index generation {}; restarting replication
from its latest generation",
+ e.generation);
+ return fetchLatestIndexOnce(forceReplication, forceCoreReload);
+ }
+ }
+
+ private IndexFetchResult fetchLatestIndexOnce(boolean forceReplication,
boolean forceCoreReload)
+ throws IOException, InterruptedException {
boolean cleanupDone = false;
boolean successfulInstall = false;
@@ -556,6 +568,7 @@ public class IndexFetcher {
log.info("Starting replication process");
// get the list of files first
fetchFileList(latestGeneration);
+ assert testWait.getAsBoolean();
// this can happen if the commit point is deleted before we fetch the
file list.
if (filesToDownload.isEmpty()) {
return IndexFetchResult.PEER_INDEX_COMMIT_DELETED;
@@ -761,6 +774,8 @@ public class IndexFetcher {
} catch (ReplicationHandlerException e) {
log.error("User aborted Replication", e);
return new
IndexFetchResult(IndexFetchResult.FAILED_BY_EXCEPTION_MESSAGE, false, e);
+ } catch (InvalidIndexGenerationException e) {
+ throw e;
} catch (SolrException e) {
throw e;
} catch (InterruptedException e) {
@@ -1593,6 +1608,8 @@ public class IndexFetcher {
bytesDownloaded = 0;
try {
fetch();
+ } catch (InvalidIndexGenerationException e) {
+ throw e;
} catch (Exception e) {
if (!aborted) {
IndexFetcher.log.error("Error fetching file, doing one retry...", e);
@@ -1605,6 +1622,7 @@ public class IndexFetcher {
}
private void fetch() throws Exception {
+ boolean invalidIndexGeneration = false;
try {
while (true) {
try (FastInputStream fis = getStream()) {
@@ -1617,17 +1635,22 @@ public class IndexFetcher {
// if there is an error continue. But continue from the point
where it got broken
}
}
+ } catch (InvalidIndexGenerationException e) {
+ invalidIndexGeneration = true;
+ throw e;
} finally {
- cleanup();
- // if cleanup succeeds, and the file is downloaded fully, then do a
fsync.
- fsyncService.execute(
- () -> {
- try {
- file.sync();
- } catch (IOException | AlreadyClosedException e) {
- fsyncException = e;
- }
- });
+ cleanup(invalidIndexGeneration);
+ if (!invalidIndexGeneration) {
+ // if cleanup succeeds, and the file is downloaded fully, then do a
fsync.
+ fsyncService.execute(
+ () -> {
+ try {
+ file.sync();
+ } catch (IOException | AlreadyClosedException e) {
+ fsyncException = e;
+ }
+ });
+ }
}
}
@@ -1744,7 +1767,7 @@ public class IndexFetcher {
}
/** cleanup everything */
- private void cleanup() {
+ private void cleanup(boolean invalidIndexGeneration) {
try {
file.close();
} catch (Exception e) {
@@ -1760,7 +1783,7 @@ public class IndexFetcher {
log.error("Error deleting file: {}", this.saveAs, e);
}
// if the failure is due to a user abort it is returned normally else
an exception is thrown
- if (!aborted)
+ if (!aborted && !invalidIndexGeneration)
throw new SolrException(
SolrException.ErrorCode.SERVER_ERROR,
"Unable to download "
@@ -1806,6 +1829,10 @@ public class IndexFetcher {
final var responseStatus = (Integer) response.get("responseStatus");
is = (InputStream) response.get("stream");
+ if (responseStatus == ErrorCode.CONFLICT.code) {
+ throw new InvalidIndexGenerationException(indexGen);
+ }
+
if (responseStatus != 200) {
final var errorMsg =
String.format(
@@ -1813,13 +1840,16 @@ public class IndexFetcher {
"Unexpected status code [%d] when downloading file [%s].",
responseStatus,
fileName);
- closeStreamAndBuildIOE(is, errorMsg, null);
+ throw closeStreamAndBuildIOE(is, errorMsg, null);
}
if (useInternalCompression) {
is = new InflaterInputStream(is);
}
return new FastInputStream(is);
+ } catch (InvalidIndexGenerationException e) {
+ IOUtils.closeQuietly(is);
+ throw e;
} catch (Exception e) {
final var ioe = closeStreamAndBuildIOE(is, "Could not download file '"
+ fileName + "'", e);
throw ioe;
@@ -1836,6 +1866,15 @@ public class IndexFetcher {
}
}
+ private static class InvalidIndexGenerationException extends IOException {
+ private final long generation;
+
+ InvalidIndexGenerationException(long generation) {
+ super("Leader no longer has index generation " + generation);
+ this.generation = generation;
+ }
+ }
+
private static class DirectoryFile implements FileInterface {
private final String saveAs;
private Directory copy2Dir;
diff --git
a/solr/core/src/java/org/apache/solr/handler/admin/api/ReplicationAPIBase.java
b/solr/core/src/java/org/apache/solr/handler/admin/api/ReplicationAPIBase.java
index 3b7170494c5..8dbf4afa911 100644
---
a/solr/core/src/java/org/apache/solr/handler/admin/api/ReplicationAPIBase.java
+++
b/solr/core/src/java/org/apache/solr/handler/admin/api/ReplicationAPIBase.java
@@ -254,6 +254,7 @@ public abstract class ReplicationAPIBase extends
JerseyResource {
protected Long indexGen;
protected IndexDeletionPolicyWrapper delPolicy;
+ private boolean commitPointSaved;
protected String fileName;
protected String cfileName;
@@ -342,7 +343,13 @@ public abstract class ReplicationAPIBase extends
JerseyResource {
// reserve commit point till write is complete
if (indexGen != null) {
- delPolicy.saveCommitPoint(indexGen);
+ try {
+ delPolicy.saveCommitPoint(indexGen);
+ commitPointSaved = true;
+ } catch (IllegalStateException e) {
+ throw new SolrException(
+ SolrException.ErrorCode.CONFLICT, "invalid index generation: " +
indexGen, e);
+ }
}
}
@@ -360,7 +367,7 @@ public abstract class ReplicationAPIBase extends
JerseyResource {
ReplicationHandler replicationHandler =
(ReplicationHandler)
solrCore.getRequestHandler(ReplicationHandler.PATH);
- if (indexGen != null) {
+ if (commitPointSaved) {
// Reserve the commit point for another 10s for the next file to be to
fetched.
// We need to keep extending the commit reservation between requests
so that the replica can
// fetch all the files correctly.
diff --git
a/solr/core/src/test/org/apache/solr/handler/TestReplicationHandler.java
b/solr/core/src/test/org/apache/solr/handler/TestReplicationHandler.java
index 326e0ac5059..cc719d054c4 100644
--- a/solr/core/src/test/org/apache/solr/handler/TestReplicationHandler.java
+++ b/solr/core/src/test/org/apache/solr/handler/TestReplicationHandler.java
@@ -34,6 +34,9 @@ import java.util.Arrays;
import java.util.Date;
import java.util.List;
import java.util.Set;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.stream.Stream;
import org.apache.lucene.index.DirectoryReader;
@@ -65,7 +68,9 @@ import org.apache.solr.common.SolrDocument;
import org.apache.solr.common.SolrDocumentList;
import org.apache.solr.common.SolrException;
import org.apache.solr.common.params.ModifiableSolrParams;
+import org.apache.solr.common.util.ExecutorUtil;
import org.apache.solr.common.util.NamedList;
+import org.apache.solr.common.util.SolrNamedThreadFactory;
import org.apache.solr.common.util.TimeSource;
import org.apache.solr.core.CachingDirectoryFactory;
import org.apache.solr.core.CoreContainer;
@@ -1493,6 +1498,63 @@ public class TestReplicationHandler extends
SolrTestCaseJ4 {
assertEquals("invalid index generation", resp.get("message"));
}
+ @Test
+ public void testFollowerRestartsWhenCommitExpiresBeforeFileDownload() throws
Exception {
+ invokeReplicationCommand(
+ buildUrl(followerJetty.getLocalPort()) + "/" + DEFAULT_TEST_CORENAME,
"disablepoll");
+
+ index(leaderClient, "id", "1", "name", "generation-g");
+ leaderClient.commit();
+ index(leaderClient, "id", "1", "name", "generation-g-plus-one");
+
+ CountDownLatch fileListFetched = new CountDownLatch(1);
+ CountDownLatch continueDownload = new CountDownLatch(1);
+ IndexFetcher.testWait =
+ () -> {
+ fileListFetched.countDown();
+ try {
+ continueDownload.await();
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new RuntimeException(e);
+ }
+ return true;
+ };
+
+ ExecutorService workload =
+ ExecutorUtil.newMDCAwareSingleThreadExecutor(
+ new SolrNamedThreadFactory("staleGenerationWorkload"));
+ try {
+ Future<?> followerFetch =
+ workload.submit(
+ () -> {
+ pullFromTo(leaderJetty, followerJetty);
+ return null;
+ });
+
+ assertTrue(fileListFetched.await(TIMEOUT, TimeUnit.MILLISECONDS));
+
+ long reserveDuration;
+ try (SolrCore core =
leaderJetty.getCoreContainer().getCore(DEFAULT_TEST_CORENAME)) {
+ ReplicationHandler handler =
+ (ReplicationHandler)
core.getRequestHandler(ReplicationHandler.PATH);
+ reserveDuration = handler.getReserveCommitDuration();
+ }
+ Thread.sleep(reserveDuration + 1000);
+ leaderClient.commit();
+
+ IndexFetcher.testWait = () -> true;
+ continueDownload.countDown();
+ followerFetch.get(TIMEOUT, TimeUnit.MILLISECONDS);
+
+ assertEquals(1, numFound(rQuery(1, "name:generation-g-plus-one",
followerClient)));
+ } finally {
+ IndexFetcher.testWait = () -> true;
+ continueDownload.countDown();
+ workload.shutdownNow();
+ }
+ }
+
@Test
public void testFetchIndexShouldReportErrorsWhenTheyOccur() throws Exception
{
int leaderPort = leaderJetty.getLocalPort();