Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 18 additions & 14 deletions src/main/java/org/opensearch/jobscheduler/sweeper/JobSweeper.java
Original file line number Diff line number Diff line change
Expand Up @@ -331,7 +331,8 @@ private void sweepAllJobIndices() {
this.lastFullSweepTimeNano = System.nanoTime();
}

private void sweepIndex(String indexName) {
@VisibleForTesting
void sweepIndex(String indexName) {
ClusterState clusterState = this.clusterService.state();
// checks to see if index no longer exists
if (!clusterState.routingTable().hasIndex(indexName)) {
Expand Down Expand Up @@ -368,14 +369,14 @@ private void sweepIndex(String indexName) {
try {
List<ShardRouting> shardRoutingList = shard.getValue();
List<String> shardNodeIds = shardRoutingList.stream().map(ShardRouting::currentNodeId).collect(Collectors.toList());
sweepShard(shard.getKey(), new ShardNodes(localNodeId, shardNodeIds), null);
sweepShard(shard.getKey(), new ShardNodes(localNodeId, shardNodeIds), -1L);
} catch (Exception e) {
log.info("Error while sweeping shard {}, error message: {}", shard.getKey(), e.getMessage());
}
}
}

private void sweepShard(ShardId shardId, ShardNodes shardNodes, String startAfter) {
private void sweepShard(ShardId shardId, ShardNodes shardNodes, long startAfter) {
ConcurrentHashMap<String, JobDocVersion> currentJobs = this.sweptJobs.containsKey(shardId)
? this.sweptJobs.get(shardId)
: new ConcurrentHashMap<>();
Expand All @@ -388,24 +389,27 @@ private void sweepShard(ShardId shardId, ShardNodes shardNodes, String startAfte
}
}

String searchAfter = startAfter == null ? "" : startAfter;
while (searchAfter != null) {
long searchAfter = startAfter;
while (searchAfter >= -1L) {
SearchRequest jobSearchRequest = new SearchRequest().indices(shardId.getIndexName())
.preference("_shards:" + shardId.id() + "|_primary")
.source(
new SearchSourceBuilder().version(true)
.seqNoAndPrimaryTerm(true)
.sort(new FieldSortBuilder("_id").unmappedType("keyword").missing("_last"))
.searchAfter(new String[] { searchAfter })
.sort(new FieldSortBuilder("_seq_no").unmappedType("long"))
.searchAfter(new Long[] { searchAfter })
.size(this.sweepPageMaxSize)
.query(QueryBuilders.matchAllQuery())
);

SearchResponse response = this.retry(
(searchRequest) -> this.client.search(searchRequest),
jobSearchRequest,
this.sweepSearchBackoff
).actionGet(this.sweepSearchTimeout);
SearchResponse response;
try {
response = this.retry((searchRequest) -> this.client.search(searchRequest), jobSearchRequest, this.sweepSearchBackoff)
.actionGet(this.sweepSearchTimeout);
} catch (Exception e) {
log.error("Aborting sweep of shard {}, will retry on next sweep cycle.", shardId, e);
return;
}
if (response.status() != RestStatus.OK) {
log.error("Error sweeping shard {}, failed querying jobs on this shard", shardId);
return;
Expand All @@ -422,10 +426,10 @@ private void sweepShard(ShardId shardId, ShardNodes shardNodes, String startAfte
}
}
if (response.getHits() == null || response.getHits().getHits().length < 1) {
searchAfter = null;
break;
} else {
SearchHit lastHit = response.getHits().getHits()[response.getHits().getHits().length - 1];
searchAfter = lastHit.getId();
searchAfter = lastHit.getSeqNo();
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import org.apache.lucene.util.BytesRef;
import org.opensearch.Version;
import org.opensearch.action.delete.DeleteResponse;
import org.opensearch.action.search.SearchResponse;
import org.opensearch.cluster.ClusterName;
import org.opensearch.cluster.ClusterState;
import org.opensearch.cluster.OpenSearchAllocationTestCase;
Expand All @@ -38,13 +39,17 @@
import org.opensearch.common.action.ActionFuture;
import org.opensearch.common.settings.ClusterSettings;
import org.opensearch.common.settings.Setting;
import org.opensearch.common.unit.TimeValue;
import org.opensearch.common.settings.Settings;
import org.opensearch.core.xcontent.NamedXContentRegistry;
import org.opensearch.core.index.Index;
import org.opensearch.index.engine.Engine;
import org.opensearch.index.mapper.ParseContext;
import org.opensearch.index.mapper.ParsedDocument;
import org.opensearch.core.index.shard.ShardId;
import org.opensearch.core.rest.RestStatus;
import org.opensearch.search.SearchHit;
import org.opensearch.search.SearchHits;
import org.opensearch.test.ClusterServiceUtils;
import org.opensearch.test.OpenSearchTestCase;
import org.opensearch.threadpool.Scheduler;
Expand Down Expand Up @@ -275,6 +280,111 @@ public void testSweep() throws IOException {
);
}

public void testSweepUsesSeqNoSort() throws IOException {
SearchHit hit = new SearchHit(1, "doc-id", null, null);
hit.sourceRef(this.getTestJsonSource());
hit.setSeqNo(42L);
hit.setPrimaryTerm(1L);
SearchHits hits = new SearchHits(new SearchHit[] { hit }, null, 1.0f);

SearchHits emptyHits = new SearchHits(new SearchHit[0], null, 1.0f);

SearchResponse firstResponse = Mockito.mock(SearchResponse.class);
Mockito.when(firstResponse.status()).thenReturn(RestStatus.OK);
Mockito.when(firstResponse.getHits()).thenReturn(hits);

SearchResponse secondResponse = Mockito.mock(SearchResponse.class);
Mockito.when(secondResponse.status()).thenReturn(RestStatus.OK);
Mockito.when(secondResponse.getHits()).thenReturn(emptyHits);

ActionFuture<SearchResponse> firstFuture = Mockito.mock(ActionFuture.class);
Mockito.when(firstFuture.actionGet(Mockito.any(TimeValue.class))).thenReturn(firstResponse);

ActionFuture<SearchResponse> secondFuture = Mockito.mock(ActionFuture.class);
Mockito.when(secondFuture.actionGet(Mockito.any(TimeValue.class))).thenReturn(secondResponse);

Mockito.when(this.client.search(Mockito.any())).thenReturn(firstFuture).thenReturn(secondFuture);

JobSweeper testSweeper = Mockito.spy(this.sweeper);
Mockito.doNothing()
.when(testSweeper)
.sweep(Mockito.any(), Mockito.anyString(), Mockito.any(BytesReference.class), Mockito.any(JobDocVersion.class));

ClusterState clusterState = buildSingleShardClusterState("index-name");
Mockito.when(this.clusterService.state()).thenReturn(clusterState);

testSweeper.sweepIndex("index-name");

// verify search was called twice: once for the page with the hit, once for the empty page
Mockito.verify(this.client, Mockito.times(2)).search(Mockito.any());
}

public void testSweepAbortsOnNonOkResponse() {
SearchResponse badResponse = Mockito.mock(SearchResponse.class);
Mockito.when(badResponse.status()).thenReturn(RestStatus.INTERNAL_SERVER_ERROR);

ActionFuture<SearchResponse> future = Mockito.mock(ActionFuture.class);
Mockito.when(future.actionGet(Mockito.any(TimeValue.class))).thenReturn(badResponse);
Mockito.when(this.client.search(Mockito.any())).thenReturn(future);

JobSweeper testSweeper = Mockito.spy(this.sweeper);
Mockito.doNothing()
.when(testSweeper)
.sweep(Mockito.any(), Mockito.anyString(), Mockito.any(BytesReference.class), Mockito.any(JobDocVersion.class));

ClusterState clusterState = buildSingleShardClusterState("index-name");
Mockito.when(this.clusterService.state()).thenReturn(clusterState);

testSweeper.sweepIndex("index-name");

// search was called once, but sweep was never called due to non-OK status
Mockito.verify(this.client, Mockito.times(1)).search(Mockito.any());
Mockito.verify(testSweeper, Mockito.times(0))
.sweep(Mockito.any(), Mockito.anyString(), Mockito.any(BytesReference.class), Mockito.any(JobDocVersion.class));
}

public void testSweepAbortsOnSearchException() {
ActionFuture<SearchResponse> failingFuture = Mockito.mock(ActionFuture.class);
Mockito.when(failingFuture.actionGet(Mockito.any(TimeValue.class)))
.thenThrow(new RuntimeException("fielddata access on _id disallowed"));
Mockito.when(this.client.search(Mockito.any())).thenReturn(failingFuture);

JobSweeper testSweeper = Mockito.spy(this.sweeper);
Mockito.doNothing()
.when(testSweeper)
.sweep(Mockito.any(), Mockito.anyString(), Mockito.any(BytesReference.class), Mockito.any(JobDocVersion.class));

ClusterState clusterState = buildSingleShardClusterState("index-name");
Mockito.when(this.clusterService.state()).thenReturn(clusterState);

// should not throw — exception is caught and logged
testSweeper.sweepIndex("index-name");

// search was attempted once before the exception aborted the loop
Mockito.verify(this.client, Mockito.times(1)).search(Mockito.any());
Mockito.verify(testSweeper, Mockito.times(0))
.sweep(Mockito.any(), Mockito.anyString(), Mockito.any(BytesReference.class), Mockito.any(JobDocVersion.class));
}

private ClusterState buildSingleShardClusterState(String indexName) {
Metadata metadata = Metadata.builder().put(createIndexMetadata(indexName, 0, 1)).build();
RoutingTable routingTable = new RoutingTable.Builder().add(
new IndexRoutingTable.Builder(metadata.index(indexName).getIndex()).initializeAsNew(metadata.index(indexName)).build()
).build();
ClusterState clusterState = ClusterState.builder(new ClusterName("cluster-name"))
.metadata(metadata)
.routingTable(routingTable)
.build();
clusterState = this.addNodesToCluter(clusterState, 1);
clusterState = this.initializeAllShards(clusterState);
// set local node so getLocalShards can match shards assigned to this node
String firstNodeId = clusterState.getNodes().iterator().next().getId();
clusterState = ClusterState.builder(clusterState)
.nodes(DiscoveryNodes.builder(clusterState.getNodes()).localNodeId(firstNodeId))
.build();
return clusterState;
}

private ClusterState addNodesToCluter(ClusterState clusterState, int nodeCount) {
DiscoveryNodes.Builder nodeBuilder = DiscoveryNodes.builder();
for (int i = 1; i <= nodeCount; i++) {
Expand Down
Loading