Do not set fatal exception when shard follow task is stopped. (#37603)
When shard follow task is cancelled while fetching operations then the fatal exception field should not be set.
This commit is contained in:
parent
614950ef1c
commit
88f4b0a326
|
@ -418,12 +418,15 @@ public abstract class ShardFollowNodeTask extends AllocatedPersistentTask {
|
||||||
|
|
||||||
private void handleFailure(Exception e, AtomicInteger retryCounter, Runnable task) {
|
private void handleFailure(Exception e, AtomicInteger retryCounter, Runnable task) {
|
||||||
assert e != null;
|
assert e != null;
|
||||||
if (shouldRetry(params.getRemoteCluster(), e) && isStopped() == false) {
|
if (shouldRetry(params.getRemoteCluster(), e)) {
|
||||||
int currentRetry = retryCounter.incrementAndGet();
|
if (isStopped() == false) {
|
||||||
LOGGER.debug(new ParameterizedMessage("{} error during follow shard task, retrying [{}]",
|
// Only retry is the shard follow task is not stopped.
|
||||||
params.getFollowShardId(), currentRetry), e);
|
int currentRetry = retryCounter.incrementAndGet();
|
||||||
long delay = computeDelay(currentRetry, params.getReadPollTimeout().getMillis());
|
LOGGER.debug(new ParameterizedMessage("{} error during follow shard task, retrying [{}]",
|
||||||
scheduler.accept(TimeValue.timeValueMillis(delay), task);
|
params.getFollowShardId(), currentRetry), e);
|
||||||
|
long delay = computeDelay(currentRetry, params.getReadPollTimeout().getMillis());
|
||||||
|
scheduler.accept(TimeValue.timeValueMillis(delay), task);
|
||||||
|
}
|
||||||
} else {
|
} else {
|
||||||
fatalException = ExceptionsHelper.convertToElastic(e);
|
fatalException = ExceptionsHelper.convertToElastic(e);
|
||||||
LOGGER.warn("shard follow task encounter non-retryable error", e);
|
LOGGER.warn("shard follow task encounter non-retryable error", e);
|
||||||
|
|
|
@ -40,7 +40,9 @@ import static org.hamcrest.Matchers.equalTo;
|
||||||
import static org.hamcrest.Matchers.greaterThanOrEqualTo;
|
import static org.hamcrest.Matchers.greaterThanOrEqualTo;
|
||||||
import static org.hamcrest.Matchers.hasSize;
|
import static org.hamcrest.Matchers.hasSize;
|
||||||
import static org.hamcrest.Matchers.instanceOf;
|
import static org.hamcrest.Matchers.instanceOf;
|
||||||
|
import static org.hamcrest.Matchers.is;
|
||||||
import static org.hamcrest.Matchers.lessThanOrEqualTo;
|
import static org.hamcrest.Matchers.lessThanOrEqualTo;
|
||||||
|
import static org.hamcrest.Matchers.nullValue;
|
||||||
import static org.hamcrest.Matchers.sameInstance;
|
import static org.hamcrest.Matchers.sameInstance;
|
||||||
|
|
||||||
public class ShardFollowNodeTaskTests extends ESTestCase {
|
public class ShardFollowNodeTaskTests extends ESTestCase {
|
||||||
|
@ -287,6 +289,35 @@ public class ShardFollowNodeTaskTests extends ESTestCase {
|
||||||
assertThat(status.leaderGlobalCheckpoint(), equalTo(63L));
|
assertThat(status.leaderGlobalCheckpoint(), equalTo(63L));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public void testFatalExceptionNotSetWhenStoppingWhileFetchingOps() {
|
||||||
|
ShardFollowTaskParams params = new ShardFollowTaskParams();
|
||||||
|
params.maxReadRequestOperationCount = 64;
|
||||||
|
params.maxOutstandingReadRequests = 1;
|
||||||
|
params.maxOutstandingWriteRequests = 1;
|
||||||
|
ShardFollowNodeTask task = createShardFollowTask(params);
|
||||||
|
startTask(task, 63, -1);
|
||||||
|
|
||||||
|
readFailures.add(new ShardNotFoundException(new ShardId("leader_index", "", 0)));
|
||||||
|
|
||||||
|
mappingVersions.add(1L);
|
||||||
|
leaderGlobalCheckpoints.add(63L);
|
||||||
|
maxSeqNos.add(63L);
|
||||||
|
responseSizes.add(64);
|
||||||
|
simulateResponse.set(true);
|
||||||
|
beforeSendShardChangesRequest = status -> {
|
||||||
|
// Cancel just before attempting to fetch operations:
|
||||||
|
task.onCancelled();
|
||||||
|
};
|
||||||
|
task.coordinateReads();
|
||||||
|
|
||||||
|
assertThat(task.isStopped(), is(true));
|
||||||
|
ShardFollowNodeTaskStatus status = task.getStatus();
|
||||||
|
assertThat(status.getFatalException(), nullValue());
|
||||||
|
assertThat(status.failedReadRequests(), equalTo(1L));
|
||||||
|
assertThat(status.successfulReadRequests(), equalTo(0L));
|
||||||
|
assertThat(status.readExceptions().size(), equalTo(1));
|
||||||
|
}
|
||||||
|
|
||||||
public void testEmptyShardChangesResponseShouldClearFetchException() {
|
public void testEmptyShardChangesResponseShouldClearFetchException() {
|
||||||
ShardFollowTaskParams params = new ShardFollowTaskParams();
|
ShardFollowTaskParams params = new ShardFollowTaskParams();
|
||||||
params.maxReadRequestOperationCount = 64;
|
params.maxReadRequestOperationCount = 64;
|
||||||
|
|
Loading…
Reference in New Issue