Skip to content
Open
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
7 changes: 7 additions & 0 deletions changelog/unreleased/SOLR-13136.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
title: Create new shards in CONSTRUCTION state so queries don't fail during shard creation
type: fixed
authors:
- name: Nick Shanin
links:
- name: SOLR-13136
url: https://issues.apache.org/jira/browse/SOLR-13136
Original file line number Diff line number Diff line change
Expand Up @@ -20,20 +20,37 @@
import static org.apache.solr.common.cloud.ZkStateReader.SHARD_ID_PROP;
import static org.apache.solr.common.params.CollectionAdminParams.FOLLOW_ALIASES;
import static org.apache.solr.common.params.CollectionParams.CollectionAction.CREATESHARD;
import static org.apache.solr.common.params.CommonAdminParams.TIMEOUT;

import java.lang.invoke.MethodHandles;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.stream.Collectors;
import org.apache.solr.cloud.ActiveReplicaWatcher;
import org.apache.solr.cloud.DistributedClusterStateUpdater;
import org.apache.solr.cloud.Overseer;
import org.apache.solr.cloud.overseer.OverseerAction;
import org.apache.solr.common.SolrCloseableLatch;
import org.apache.solr.common.SolrException;
import org.apache.solr.common.cloud.ClusterState;
import org.apache.solr.common.cloud.DocCollection;
import org.apache.solr.common.cloud.Replica;
import org.apache.solr.common.cloud.ReplicaCount;
import org.apache.solr.common.cloud.Slice;
import org.apache.solr.common.cloud.ZkNodeProps;
import org.apache.solr.common.cloud.ZkStateReader;
import org.apache.solr.common.params.CollectionParams;
import org.apache.solr.common.params.CommonAdminParams;
import org.apache.solr.common.params.CoreAdminParams;
import org.apache.solr.common.params.ModifiableSolrParams;
import org.apache.solr.common.util.NamedList;
import org.apache.solr.common.util.SimpleOrderedMap;
import org.apache.solr.common.util.Utils;
import org.apache.solr.handler.component.ShardHandler;
import org.apache.zookeeper.KeeperException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -67,6 +84,7 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList
}
ClusterState clusterState = adminCmdContext.getClusterState();
DocCollection collection = clusterState.getCollection(collectionName);
boolean sliceAlreadyExists = collection.getSlice(sliceName) != null;

ReplicaCount numReplicas = ReplicaCount.fromMessage(message, collection, 1);
if (!numReplicas.hasLeaderReplica()) {
Expand Down Expand Up @@ -115,6 +133,7 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList

CollectionHandlingUtils.addPropertyParams(message, addReplicasProps);
final NamedList<Object> addResult = new NamedList<>();
int timeout = message.getInt(TIMEOUT, 10 * 60); // 10 minutes
try {
new AddReplicaCmd(ccc)
.addReplica(
Expand Down Expand Up @@ -148,18 +167,162 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList
success.addAll(addResultSuccess);
}
});
} catch (Assign.AssignmentException e) {
// clean up the slice that we created
new DeleteShardCmd(ccc)
.call(
adminCmdContext
.subRequestContext(CollectionParams.CollectionAction.DELETESHARD, async)
.withClusterState(clusterState),
new ZkNodeProps(COLLECTION_PROP, collectionName, SHARD_ID_PROP, sliceName),
results);
} catch (Exception e) {
if (!sliceAlreadyExists) {
// Don't leave a half-created shard stuck in CONSTRUCTION; remove it so the create can be
// retried cleanly. DeleteShardCmd needs a fresh cluster state here: the snapshot from
// waitForNewShard predates AddReplicaCmd, so it shows a slice with no replicas and the
// delete would leave the added replicas' cores behind.
try {
new DeleteShardCmd(ccc)
.call(
adminCmdContext
.subRequestContext(CollectionParams.CollectionAction.DELETESHARD, async)
.withClusterState(ccc.getZkStateReader().getClusterState()),
new ZkNodeProps(COLLECTION_PROP, collectionName, SHARD_ID_PROP, sliceName),
results);
} catch (Exception cleanupEx) {
log.warn("Failed to delete shard {} after failed create", sliceName, cleanupEx);
}
}
throw e;
}

if (!sliceAlreadyExists) {
// The new slice is in CONSTRUCTION state, so queries skip it until it is activated. Only
// activate it once its replicas are active and any buffered updates are applied. If that
// fails, delete the half-created slice instead of activating a shard that may not be
// ready, mirroring the AddReplica failure handling above, and report the failure so the
// create can be retried.
try {
waitForShardReplicasActive(collectionName, sliceName, numReplicas.total(), timeout);
applyBufferedUpdatesOnLeader(adminCmdContext, collectionName, sliceName);
} catch (Exception e) {
try {
new DeleteShardCmd(ccc)
.call(
adminCmdContext
.subRequestContext(CollectionParams.CollectionAction.DELETESHARD, async)
.withClusterState(ccc.getZkStateReader().getClusterState()),
new ZkNodeProps(COLLECTION_PROP, collectionName, SHARD_ID_PROP, sliceName),
results);
} catch (Exception cleanupEx) {
log.warn("Failed to delete shard {} after failed create", sliceName, cleanupEx);
}
throw e;
}
activateShard(collectionName, sliceName);
}

log.info("Finished create command on all shards for collection: {}", collectionName);
}

/**
* Cores of a CONSTRUCTION shard buffer their updates, which a split lifts for its sub-shards but
* nothing else would for a created shard, so ask the leader to apply them.
*/
private void applyBufferedUpdatesOnLeader(
AdminCmdContext adminCmdContext, String collectionName, String sliceName) {
Slice slice =
ccc.getZkStateReader().getClusterState().getCollection(collectionName).getSlice(sliceName);
Replica leader = slice == null ? null : slice.getLeader();
if (leader == null) {
log.warn("No leader for new shard {} of collection {}", sliceName, collectionName);
return;
}
ModifiableSolrParams params = new ModifiableSolrParams();
params.set(
CoreAdminParams.ACTION, CoreAdminParams.CoreAdminAction.REQUESTAPPLYUPDATES.toString());
params.set(CoreAdminParams.NAME, leader.getCoreName());

ShardHandler shardHandler = ccc.newShardHandler();
CollectionHandlingUtils.ShardRequestTracker tracker =
CollectionHandlingUtils.asyncRequestTracker(adminCmdContext, ccc);
tracker.sendShardRequest(leader, params, shardHandler);

// The core refuses the request when it is not buffering, which is fine, so only log a failure.
NamedList<Object> applyResults = new NamedList<>();
tracker.processResponses(applyResults, shardHandler, false, null);
if (applyResults.get("failure") != null) {
log.warn(
"Leader {} of new shard {} did not apply buffered updates: {}",
leader.getName(),
sliceName,
applyResults.get("failure"));
}
}

/** Flip a newly created shard from CONSTRUCTION to ACTIVE so it becomes visible to queries. */
private void activateShard(String collectionName, String sliceName)
throws KeeperException, InterruptedException {
Map<String, Object> activateProps = new HashMap<>();
activateProps.put(Overseer.QUEUE_OPERATION, OverseerAction.UPDATESHARDSTATE.toLower());
activateProps.put(COLLECTION_PROP, collectionName);
activateProps.put(sliceName, Slice.State.ACTIVE.toString());
ZkNodeProps activateMsg = new ZkNodeProps(activateProps);
if (ccc.getDistributedClusterStateUpdater().isDistributedStateUpdate()) {
ccc.getDistributedClusterStateUpdater()
.doSingleStateUpdate(
DistributedClusterStateUpdater.MutatingCommand.SliceUpdateShardState,
activateMsg,
ccc.getSolrCloudManager(),
ccc.getZkStateReader());
} else {
ccc.offerStateUpdate(activateMsg);
}
}

/** Wait until a shard has the expected number of replicas and all of them are ACTIVE. */
private void waitForShardReplicasActive(
String collectionName, String sliceName, int expectedReplicas, int timeout)
throws InterruptedException {
ZkStateReader zkStateReader = ccc.getZkStateReader();
try {
zkStateReader.waitForState(
collectionName,
timeout,
TimeUnit.SECONDS,
collection -> {
Slice newSlice = collection == null ? null : collection.getSlice(sliceName);
return newSlice != null && newSlice.getReplicas().size() >= expectedReplicas;
});
} catch (TimeoutException e) {
throw new SolrException(
SolrException.ErrorCode.SERVER_ERROR,
"Timeout waiting " + timeout + " seconds for the replicas of shard " + sliceName,
e);
}
Slice slice = zkStateReader.getClusterState().getCollection(collectionName).getSlice(sliceName);
if (slice == null) {
throw new SolrException(
SolrException.ErrorCode.SERVER_ERROR,
"Newly created shard " + sliceName + " not visible in cluster state");
}
List<String> coreNames =
slice.getReplicas().stream().map(Replica::getCoreName).collect(Collectors.toList());
if (coreNames.isEmpty()) {
log.warn(
"No replicas found for new shard {} of collection {}; activating without waiting",
sliceName,
collectionName);
return;
}
SolrCloseableLatch latch =
new SolrCloseableLatch(coreNames.size(), ccc.getCloseableToLatchOn());
ActiveReplicaWatcher watcher = new ActiveReplicaWatcher(collectionName, null, coreNames, latch);
try {
zkStateReader.registerCollectionStateWatcher(collectionName, watcher);
if (!latch.await(timeout, TimeUnit.SECONDS)) {
throw new SolrException(
SolrException.ErrorCode.SERVER_ERROR,
"Timeout waiting "
+ timeout
+ " seconds for replicas of shard "
+ sliceName
+ " to become active.");
}
} finally {
zkStateReader.removeCollectionStateWatcher(collectionName, watcher);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,10 @@ public ZkWriteCommand createShard(final ClusterState clusterState, ZkNodeProps m
Map<String, Object> sliceProps = new HashMap<>();
String shardRange = message.getStr(ZkStateReader.SHARD_RANGE_PROP);
String shardState = message.getStr(ZkStateReader.SHARD_STATE_PROP);
if (shardState == null) {
// A brand-new slice has no replicas yet; keep it out of query routing until they are added
shardState = Slice.State.CONSTRUCTION.toString();
}
String shardParent = message.getStr(ZkStateReader.SHARD_PARENT_PROP);
String shardParentZkSession = message.getStr("shard_parent_zk_session");
String shardParentNode = message.getStr("shard_parent_node");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,7 @@ void setShard(String shard) {
this.shard = shard;
}

void setException(Throwable exception) {
public void setException(Throwable exception) {
this.exception = exception;
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,132 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.solr.cloud;

import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.solr.client.solrj.request.CollectionAdminRequest;
import org.apache.solr.client.solrj.request.UpdateRequest;
import org.apache.solr.common.cloud.CollectionStateWatcher;
import org.apache.solr.common.cloud.Replica;
import org.apache.solr.common.cloud.Slice;
import org.apache.solr.common.cloud.ZkStateReader;
import org.apache.solr.common.params.ModifiableSolrParams;
import org.apache.solr.common.params.ShardParams;
import org.junit.BeforeClass;
import org.junit.Test;

/**
* Verifies that a shard created via the Collections API stays in CONSTRUCTION state until its
* replicas are ACTIVE, so queries never route to a half-created shard (SOLR-13136).
*/
public class CreateShardConstructionTest extends SolrCloudTestCase {

@BeforeClass
public static void setupCluster() throws Exception {
configureCluster(2).addConfig("conf", configset("cloud-minimal")).configure();
}

@Test
public void testNewShardStaysConstructionUntilReplicasActive() throws Exception {
String collection = "shard_construction_test";
CollectionAdminRequest.createCollectionWithImplicitRouter(collection, "conf", "shard1", 1)
.processAndWait(cluster.getSolrClient(), 60);
cluster.waitForActiveCollection(collection, 1, 1);

// Record every state the new shard is observed in; a ZK watcher fires on each state change
// so the brief CONSTRUCTION window can't be missed by polling.
Set<Slice.State> observedStates = ConcurrentHashMap.newKeySet();
CountDownLatch activeLatch = new CountDownLatch(1);
CollectionStateWatcher watcher =
(liveNodes, collectionState) -> {
if (collectionState != null) {
Slice slice = collectionState.getSlice("shard2");
if (slice != null) {
observedStates.add(slice.getState());
if (slice.getState() == Slice.State.ACTIVE
&& slice.getReplicas().stream()
.allMatch(r -> r.getState() == Replica.State.ACTIVE)) {
activeLatch.countDown();
}
}
}
return false;
};
ZkStateReader zkStateReader = ZkStateReader.from(cluster.getSolrClient());
zkStateReader.registerCollectionStateWatcher(collection, watcher);
try {
AtomicReference<Exception> createError = new AtomicReference<>();
Thread createThread =
new Thread(
() -> {
try {
CollectionAdminRequest.createShard(collection, "shard2")
.process(cluster.getSolrClient());
} catch (Exception e) {
createError.set(e);
}
});
createThread.start();

// Query while the shard is under construction; it must not fail with
// "no servers hosting shard".
ModifiableSolrParams queryParams = params("q", "*:*", "rows", "0");
Exception queryError = null;
long deadlineNanos = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(120000);
while (!activeLatch.await(200, TimeUnit.MILLISECONDS)) {
if (System.nanoTime() > deadlineNanos) {
break;
}
try {
cluster.getSolrClient().query(collection, queryParams);
} catch (Exception e) {
if (e.getMessage() != null && e.getMessage().contains("no servers hosting shard")) {
queryError = e;
break;
}
}
}
createThread.join(180000);

assertNull("createShard failed: " + createError.get(), createError.get());
assertNull("query failed during shard creation: " + queryError, queryError);
assertTrue(
"new shard was never observed in CONSTRUCTION state (saw: " + observedStates + ")",
observedStates.contains(Slice.State.CONSTRUCTION));
assertTrue("new shard never reached ACTIVE", activeLatch.getCount() == 0);

// updates to the new shard must be applied and searchable once it is active
UpdateRequest update = new UpdateRequest();
update.add("id", "doc1");
update.setParam(ShardParams._ROUTE_, "shard2");
update.commit(cluster.getSolrClient(), collection);
assertEquals(
1,
cluster
.getSolrClient()
.query(collection, params("q", "*:*", ShardParams._ROUTE_, "shard2"))
.getResults()
.getNumFound());
} finally {
zkStateReader.removeCollectionStateWatcher(collection, watcher);
CollectionAdminRequest.deleteCollection(collection).process(cluster.getSolrClient());
}
}
}
Loading
Loading