forked from opensearch-project/OpenSearch
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Added transport action for bulk async shard fetch for primary shards (o…
…pensearch-project#8218) * Added transport action for bulk async shard fetch for primary shards Signed-off-by: sudarshan baliga <baliga108@gmail.com> Co-authored-by: Shivansh Arora <hishiv@amazon.com>
- Loading branch information
Showing
6 changed files
with
708 additions
and
70 deletions.
There are no files selected for viewing
77 changes: 77 additions & 0 deletions
77
server/src/internalClusterTest/java/org/opensearch/gateway/GatewayRecoveryTestUtils.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,77 @@ | ||
/* | ||
* SPDX-License-Identifier: Apache-2.0 | ||
* | ||
* The OpenSearch Contributors require contributions made to | ||
* this file be licensed under the Apache-2.0 license or a | ||
* compatible open source license. | ||
*/ | ||
|
||
package org.opensearch.gateway; | ||
|
||
import org.opensearch.action.admin.cluster.state.ClusterStateRequest; | ||
import org.opensearch.action.admin.cluster.state.ClusterStateResponse; | ||
import org.opensearch.cluster.metadata.IndexMetadata; | ||
import org.opensearch.cluster.node.DiscoveryNode; | ||
import org.opensearch.core.index.Index; | ||
import org.opensearch.core.index.shard.ShardId; | ||
import org.opensearch.env.NodeEnvironment; | ||
import org.opensearch.index.shard.ShardPath; | ||
import org.opensearch.indices.store.ShardAttributes; | ||
|
||
import java.io.IOException; | ||
import java.nio.file.DirectoryStream; | ||
import java.nio.file.Files; | ||
import java.nio.file.Path; | ||
import java.util.HashMap; | ||
import java.util.LinkedList; | ||
import java.util.List; | ||
import java.util.Map; | ||
import java.util.concurrent.ExecutionException; | ||
|
||
import static org.opensearch.test.OpenSearchIntegTestCase.client; | ||
import static org.opensearch.test.OpenSearchIntegTestCase.internalCluster; | ||
import static org.opensearch.test.OpenSearchIntegTestCase.resolveIndex; | ||
|
||
public class GatewayRecoveryTestUtils { | ||
|
||
public static DiscoveryNode[] getDiscoveryNodes() throws ExecutionException, InterruptedException { | ||
final ClusterStateRequest clusterStateRequest = new ClusterStateRequest(); | ||
clusterStateRequest.local(false); | ||
clusterStateRequest.clear().nodes(true).routingTable(true).indices("*"); | ||
ClusterStateResponse clusterStateResponse = client().admin().cluster().state(clusterStateRequest).get(); | ||
final List<DiscoveryNode> nodes = new LinkedList<>(clusterStateResponse.getState().nodes().getDataNodes().values()); | ||
DiscoveryNode[] disNodesArr = new DiscoveryNode[nodes.size()]; | ||
nodes.toArray(disNodesArr); | ||
return disNodesArr; | ||
} | ||
|
||
public static Map<ShardId, ShardAttributes> prepareRequestMap(String[] indices, int primaryShardCount) { | ||
Map<ShardId, ShardAttributes> shardIdShardAttributesMap = new HashMap<>(); | ||
for (String indexName : indices) { | ||
final Index index = resolveIndex(indexName); | ||
final String customDataPath = IndexMetadata.INDEX_DATA_PATH_SETTING.get( | ||
client().admin().indices().prepareGetSettings(indexName).get().getIndexToSettings().get(indexName) | ||
); | ||
for (int shardIdNum = 0; shardIdNum < primaryShardCount; shardIdNum++) { | ||
final ShardId shardId = new ShardId(index, shardIdNum); | ||
shardIdShardAttributesMap.put(shardId, new ShardAttributes(shardId, customDataPath)); | ||
} | ||
} | ||
return shardIdShardAttributesMap; | ||
} | ||
|
||
public static void corruptShard(String nodeName, ShardId shardId) throws IOException, InterruptedException { | ||
for (Path path : internalCluster().getInstance(NodeEnvironment.class, nodeName).availableShardPaths(shardId)) { | ||
final Path indexPath = path.resolve(ShardPath.INDEX_FOLDER_NAME); | ||
if (Files.exists(indexPath)) { // multi data path might only have one path in use | ||
try (DirectoryStream<Path> stream = Files.newDirectoryStream(indexPath)) { | ||
for (Path item : stream) { | ||
if (item.getFileName().toString().startsWith("segments_")) { | ||
Files.delete(item); | ||
} | ||
} | ||
} | ||
} | ||
} | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
114 changes: 114 additions & 0 deletions
114
server/src/main/java/org/opensearch/gateway/TransportNodesGatewayStartedShardHelper.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,114 @@ | ||
/* | ||
* SPDX-License-Identifier: Apache-2.0 | ||
* | ||
* The OpenSearch Contributors require contributions made to | ||
* this file be licensed under the Apache-2.0 license or a | ||
* compatible open source license. | ||
*/ | ||
|
||
package org.opensearch.gateway; | ||
|
||
import org.apache.logging.log4j.Logger; | ||
import org.apache.logging.log4j.message.ParameterizedMessage; | ||
import org.opensearch.OpenSearchException; | ||
import org.opensearch.cluster.metadata.IndexMetadata; | ||
import org.opensearch.cluster.service.ClusterService; | ||
import org.opensearch.common.settings.Settings; | ||
import org.opensearch.core.index.shard.ShardId; | ||
import org.opensearch.core.xcontent.NamedXContentRegistry; | ||
import org.opensearch.env.NodeEnvironment; | ||
import org.opensearch.index.IndexSettings; | ||
import org.opensearch.index.shard.IndexShard; | ||
import org.opensearch.index.shard.ShardPath; | ||
import org.opensearch.index.shard.ShardStateMetadata; | ||
import org.opensearch.index.store.Store; | ||
import org.opensearch.indices.IndicesService; | ||
|
||
import java.io.IOException; | ||
|
||
/** | ||
* This class has the common code used in {@link TransportNodesListGatewayStartedShards} and | ||
* {@link TransportNodesListGatewayStartedShardsBatch} to get the shard info on the local node. | ||
* <p> | ||
* This class should not be used to add more functions and will be removed when the | ||
* {@link TransportNodesListGatewayStartedShards} will be deprecated and all the code will be moved to | ||
* {@link TransportNodesListGatewayStartedShardsBatch} | ||
* | ||
* @opensearch.internal | ||
*/ | ||
public class TransportNodesGatewayStartedShardHelper { | ||
public static TransportNodesListGatewayStartedShardsBatch.NodeGatewayStartedShard getShardInfoOnLocalNode( | ||
Logger logger, | ||
final ShardId shardId, | ||
NamedXContentRegistry namedXContentRegistry, | ||
NodeEnvironment nodeEnv, | ||
IndicesService indicesService, | ||
String shardDataPathInRequest, | ||
Settings settings, | ||
ClusterService clusterService | ||
) throws IOException { | ||
logger.trace("{} loading local shard state info", shardId); | ||
ShardStateMetadata shardStateMetadata = ShardStateMetadata.FORMAT.loadLatestState( | ||
logger, | ||
namedXContentRegistry, | ||
nodeEnv.availableShardPaths(shardId) | ||
); | ||
if (shardStateMetadata != null) { | ||
if (indicesService.getShardOrNull(shardId) == null | ||
&& shardStateMetadata.indexDataLocation == ShardStateMetadata.IndexDataLocation.LOCAL) { | ||
final String customDataPath; | ||
if (shardDataPathInRequest != null) { | ||
customDataPath = shardDataPathInRequest; | ||
} else { | ||
// TODO: Fallback for BWC with older OpenSearch versions. | ||
// Remove once request.getCustomDataPath() always returns non-null | ||
final IndexMetadata metadata = clusterService.state().metadata().index(shardId.getIndex()); | ||
if (metadata != null) { | ||
customDataPath = new IndexSettings(metadata, settings).customDataPath(); | ||
} else { | ||
logger.trace("{} node doesn't have meta data for the requests index", shardId); | ||
throw new OpenSearchException("node doesn't have meta data for index " + shardId.getIndex()); | ||
} | ||
} | ||
// we don't have an open shard on the store, validate the files on disk are openable | ||
ShardPath shardPath = null; | ||
try { | ||
shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, customDataPath); | ||
if (shardPath == null) { | ||
throw new IllegalStateException(shardId + " no shard path found"); | ||
} | ||
Store.tryOpenIndex(shardPath.resolveIndex(), shardId, nodeEnv::shardLock, logger); | ||
} catch (Exception exception) { | ||
final ShardPath finalShardPath = shardPath; | ||
logger.trace( | ||
() -> new ParameterizedMessage( | ||
"{} can't open index for shard [{}] in path [{}]", | ||
shardId, | ||
shardStateMetadata, | ||
(finalShardPath != null) ? finalShardPath.resolveIndex() : "" | ||
), | ||
exception | ||
); | ||
String allocationId = shardStateMetadata.allocationId != null ? shardStateMetadata.allocationId.getId() : null; | ||
return new TransportNodesListGatewayStartedShardsBatch.NodeGatewayStartedShard( | ||
allocationId, | ||
shardStateMetadata.primary, | ||
null, | ||
exception | ||
); | ||
} | ||
} | ||
|
||
logger.debug("{} shard state info found: [{}]", shardId, shardStateMetadata); | ||
String allocationId = shardStateMetadata.allocationId != null ? shardStateMetadata.allocationId.getId() : null; | ||
final IndexShard shard = indicesService.getShardOrNull(shardId); | ||
return new TransportNodesListGatewayStartedShardsBatch.NodeGatewayStartedShard( | ||
allocationId, | ||
shardStateMetadata.primary, | ||
shard != null ? shard.getLatestReplicationCheckpoint() : null | ||
); | ||
} | ||
logger.trace("{} no local shard info found", shardId); | ||
return new TransportNodesListGatewayStartedShardsBatch.NodeGatewayStartedShard(null, false, null); | ||
} | ||
} |
Oops, something went wrong.