From a3a798d7fd9dd87c6c757b845d417602dd457cbf Mon Sep 17 00:00:00 2001 From: pzurek Date: Wed, 9 Mar 2022 15:32:20 +0100 Subject: [PATCH] MNT-22770 Fix duplicates for sharded reindexing --- .../solr/tracker/MetadataTracker.java | 99 +++++----- ...frescoSolrMetadataTrackerReindexingIT.java | 183 ++++++++++++++++++ 2 files changed, 228 insertions(+), 54 deletions(-) create mode 100644 search-services/alfresco-search/src/test/java/org/alfresco/solr/tracker/DistributedAlfrescoSolrMetadataTrackerReindexingIT.java diff --git a/search-services/alfresco-search/src/main/java/org/alfresco/solr/tracker/MetadataTracker.java b/search-services/alfresco-search/src/main/java/org/alfresco/solr/tracker/MetadataTracker.java index 8e9d002ca..be14895b0 100644 --- a/search-services/alfresco-search/src/main/java/org/alfresco/solr/tracker/MetadataTracker.java +++ b/search-services/alfresco-search/src/main/java/org/alfresco/solr/tracker/MetadataTracker.java @@ -2,7 +2,7 @@ * #%L * Alfresco Search Services * %% - * Copyright (C) 2005 - 2020 Alfresco Software Limited + * Copyright (C) 2005 - 2022 Alfresco Software Limited * %% * This file is part of the Alfresco software. * If the software was purchased under a paid Alfresco license, the terms of @@ -594,32 +594,24 @@ public class MetadataTracker extends ActivatableTracker private void reindexNodes() throws IOException, AuthenticationException, JSONException { - boolean requiresCommit = false; while (nodesToReindex.peek() != null) { Long nodeId = nodesToReindex.poll(); if (nodeId != null) { - // make sure it is cleaned out so we do not miss deletes - this.infoSrv.deleteByNodeId(nodeId); - - Node node = new Node(); + final Node node = new Node(); node.setId(nodeId); node.setStatus(SolrApiNodeStatus.UNKNOWN); node.setTxnId(Long.MAX_VALUE); - this.infoSrv.indexNode(node, true); - LOGGER.info("REINDEX ACTION - Node {} has been reindexed", node.getId()); - requiresCommit = true; + for (final Node n : filterNodes(List.of(node))) + { + this.infoSrv.indexNode(n, true); + LOGGER.info("REINDEX ACTION - Node {} has been reindexed", n.getId()); + } } checkShutdown(); } - - if(requiresCommit) - { - checkShutdown(); - //this.infoSrv.commit(); - } } private void reindexNodesByQuery() throws IOException, AuthenticationException, JSONException @@ -644,6 +636,44 @@ public class MetadataTracker extends ActivatableTracker } } + private List filterNodes(List nodes) + { + List filteredList = new ArrayList<>(nodes.size()); + for(Node node : nodes) + { + if(docRouter.routeNode(shardCount, shardInstance, node)) + { + filteredList.add(node); + } + else if (cascadeTrackerEnabled) + { + if(node.getStatus() == SolrApiNodeStatus.UPDATED) + { + Node doCascade = new Node(); + doCascade.setAclId(node.getAclId()); + doCascade.setId(node.getId()); + doCascade.setNodeRef(node.getNodeRef()); + doCascade.setStatus(SolrApiNodeStatus.NON_SHARD_UPDATED); + doCascade.setTenant(node.getTenant()); + doCascade.setTxnId(node.getTxnId()); + filteredList.add(doCascade); + } + else // DELETED & UNKNOWN + { + // Make sure anything no longer relevant to this shard is deleted. + Node doDelete = new Node(); + doDelete.setAclId(node.getAclId()); + doDelete.setId(node.getId()); + doDelete.setNodeRef(node.getNodeRef()); + doDelete.setStatus(SolrApiNodeStatus.NON_SHARD_DELETED); + doDelete.setTenant(node.getTenant()); + doDelete.setTxnId(node.getTxnId()); + filteredList.add(doDelete); + } + } + } + return filteredList; + } private void purgeTransactions() throws IOException, JSONException { @@ -1164,45 +1194,6 @@ public class MetadataTracker extends ActivatableTracker { setRollback(true, failCausedBy); } - - private List filterNodes(List nodes) - { - List filteredList = new ArrayList<>(nodes.size()); - for(Node node : nodes) - { - if(docRouter.routeNode(shardCount, shardInstance, node)) - { - filteredList.add(node); - } - else if (cascadeTrackerEnabled) - { - if(node.getStatus() == SolrApiNodeStatus.UPDATED) - { - Node doCascade = new Node(); - doCascade.setAclId(node.getAclId()); - doCascade.setId(node.getId()); - doCascade.setNodeRef(node.getNodeRef()); - doCascade.setStatus(SolrApiNodeStatus.NON_SHARD_UPDATED); - doCascade.setTenant(node.getTenant()); - doCascade.setTxnId(node.getTxnId()); - filteredList.add(doCascade); - } - else // DELETED & UNKNOWN - { - // Make sure anything no longer relevant to this shard is deleted. - Node doDelete = new Node(); - doDelete.setAclId(node.getAclId()); - doDelete.setId(node.getId()); - doDelete.setNodeRef(node.getNodeRef()); - doDelete.setStatus(SolrApiNodeStatus.NON_SHARD_DELETED); - doDelete.setTenant(node.getTenant()); - doDelete.setTxnId(node.getTxnId()); - filteredList.add(doDelete); - } - } - } - return filteredList; - } } @Override diff --git a/search-services/alfresco-search/src/test/java/org/alfresco/solr/tracker/DistributedAlfrescoSolrMetadataTrackerReindexingIT.java b/search-services/alfresco-search/src/test/java/org/alfresco/solr/tracker/DistributedAlfrescoSolrMetadataTrackerReindexingIT.java new file mode 100644 index 000000000..c5ee0172c --- /dev/null +++ b/search-services/alfresco-search/src/test/java/org/alfresco/solr/tracker/DistributedAlfrescoSolrMetadataTrackerReindexingIT.java @@ -0,0 +1,183 @@ +/* + * #%L + * Alfresco Search Services + * %% + * Copyright (C) 2005 - 2022 Alfresco Software Limited + * %% + * This file is part of the Alfresco software. + * If the software was purchased under a paid Alfresco license, the terms of + * the paid license agreement will prevail. Otherwise, the software is + * provided under the following open source license terms: + * + * Alfresco is free software: you can redistribute it and/or modify + * it under the terms of the GNU Lesser General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * Alfresco is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Lesser General Public License for more details. + * + * You should have received a copy of the GNU Lesser General Public License + * along with Alfresco. If not, see . + * #L% + */ +package org.alfresco.solr.tracker; + +import static org.alfresco.solr.AlfrescoSolrUtils.getNode; +import static org.alfresco.solr.AlfrescoSolrUtils.getNodeMetaData; +import static org.alfresco.solr.AlfrescoSolrUtils.getTransaction; + +import java.io.IOException; +import java.time.Duration; +import java.time.Instant; +import java.util.List; +import java.util.Map; +import java.util.function.Function; +import java.util.stream.Collectors; + +import org.alfresco.repo.search.adaptor.QueryConstants; +import org.alfresco.solr.AbstractAlfrescoDistributedIT; +import org.alfresco.solr.AlfrescoCoreAdminHandler; +import org.alfresco.solr.AlfrescoSolrUtils.TestActChanges; +import org.alfresco.solr.client.Acl; +import org.alfresco.solr.client.AclChangeSet; +import org.alfresco.solr.client.Node; +import org.alfresco.solr.client.NodeMetaData; +import org.alfresco.solr.client.Transaction; +import org.apache.lucene.index.Term; +import org.apache.lucene.search.BooleanClause; +import org.apache.lucene.search.BooleanQuery; +import org.apache.lucene.search.LegacyNumericRangeQuery; +import org.apache.lucene.search.TermQuery; +import org.apache.solr.SolrTestCaseJ4; +import org.apache.solr.client.solrj.SolrClient; +import org.apache.solr.client.solrj.SolrServerException; +import org.apache.solr.client.solrj.response.QueryResponse; +import org.apache.solr.common.SolrDocumentList; +import org.apache.solr.common.params.CoreAdminParams; +import org.apache.solr.common.params.ModifiableSolrParams; +import org.apache.solr.core.SolrCore; +import org.apache.solr.request.LocalSolrQueryRequest; +import org.apache.solr.request.SolrQueryRequest; +import org.apache.solr.response.SolrQueryResponse; +import org.junit.AfterClass; +import org.junit.BeforeClass; +import org.junit.Test; + +@SolrTestCaseJ4.SuppressSSL +public class DistributedAlfrescoSolrMetadataTrackerReindexingIT extends AbstractAlfrescoDistributedIT +{ + private static final int NUMBER_OF_SHARDS = 2; + private static final Duration TEST_DURATION = Duration.ofMinutes(1); + + @BeforeClass + public static void initData() throws Throwable + { + initSolrServers(NUMBER_OF_SHARDS, getSimpleClassName(), null); + } + + @AfterClass + public static void destroyData() + { + dismissSolrServers(); + } + + @Test + public void shouldReindexNodeOnlyOnShardWhereNodeBelongsTo() throws Exception + { + putHandleDefaults(); + + long nodeId = indexNewNode(); + + assertThatNodeIsIndexedOnExactlyOneShard(nodeId); + + final Instant end = Instant.now().plus(TEST_DURATION); + while (Instant.now().isBefore(end)) + { + triggerNodeIdReindexingOnAllShards(nodeId); + assertThatNodeIsIndexedOnExactlyOneShard(nodeId); + Thread.sleep(7_000); + } + } + + private void triggerNodeIdReindexingOnAllShards(long nodeId) throws Exception + { + final Map cores = getCores(solrShards) + .stream() + .collect(Collectors.toMap(SolrCore::getName, Function.identity())); + + final List adminHandlers = getAdminHandlers(solrShards); + assertEquals(NUMBER_OF_SHARDS, adminHandlers.size()); + + for (int i = 0; i < NUMBER_OF_SHARDS; i++) + { + final SolrCore core = cores.get("shard" + i); + assertNotNull(core); + + final AlfrescoCoreAdminHandler admin = adminHandlers.get(i); + assertNotNull(admin); + + final SolrQueryRequest reindexRequest = new LocalSolrQueryRequest(core, + params( + CoreAdminParams.ACTION, "REINDEX", + CoreAdminParams.CORE, core.getName(), + "nodeid", Long.toString(nodeId) + ) + ); + final SolrQueryResponse reindexResponse = new SolrQueryResponse(); + + admin.handleRequestBody(reindexRequest, reindexResponse); + assertNull(reindexResponse.getException()); + } + } + + private void assertThatNodeIsIndexedOnExactlyOneShard(long nodeId) + { + final List allClients = getShardedClients(); + assertEquals(NUMBER_OF_SHARDS, allClients.size()); + + final ModifiableSolrParams queryParams = params("qt", "/afts", "q", "DBID:" + nodeId); + + final long sumForAllShards = allClients.stream().map(c -> { + try + { + return c.query(queryParams); + } catch (SolrServerException | IOException e) + { + throw new RuntimeException(e); + } + }).map(QueryResponse::getResults).mapToLong(SolrDocumentList::getNumFound).sum(); + + assertEquals(1, sumForAllShards); + } + + private long indexNewNode() throws Exception + { + final Transaction tx = getTransaction(0, 1); + final Acl acl = getAcl(); + + final Node node = getNode(tx, acl, Node.SolrApiNodeStatus.UPDATED); + final NodeMetaData nodeMetaData = getNodeMetaData(node, tx, acl, "piotrek", null, false); + + indexTransaction(tx, List.of(node), List.of(nodeMetaData)); + waitForDocCount(new TermQuery(new Term("content@s___t@{http://www.alfresco.org/model/content/1.0}content", "world")), 1, 100000); + + return node.getId(); + } + + private Acl getAcl() throws Exception + { + TestActChanges testActChanges = new TestActChanges().createBasicTestData(); + AclChangeSet aclChangeSet = testActChanges.getChangeSet(); + + BooleanQuery.Builder builder = new BooleanQuery.Builder(); + builder.add(new BooleanClause(new TermQuery(new Term(QueryConstants.FIELD_SOLR4_ID, "TRACKER!STATE!ACLTX")), BooleanClause.Occur.MUST)); + builder.add(new BooleanClause(LegacyNumericRangeQuery.newLongRange(QueryConstants.FIELD_S_ACLTXID, aclChangeSet.getId(), aclChangeSet.getId() + 1, true, false), BooleanClause.Occur.MUST)); + BooleanQuery waitForQuery = builder.build(); + waitForDocCountAllCores(waitForQuery, 1, 80000); + + return testActChanges.getFirstAcl(); + } +}