Loading shared/src/main/java/com/inteligr8/alfresco/asie/rest/AbstractActionWebScript.java +0 −6 Original line number Diff line number Diff line Loading @@ -9,14 +9,8 @@ import java.util.Map; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.alfresco.model.ContentModel; import org.alfresco.service.cmr.repository.InvalidNodeRefException; import org.alfresco.service.cmr.repository.NodeRef; import org.alfresco.service.cmr.repository.NodeService; import org.alfresco.service.cmr.repository.StoreRef; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.extensions.webscripts.WebScriptException; import org.springframework.extensions.webscripts.WebScriptRequest; import org.springframework.extensions.webscripts.WebScriptResponse; Loading shared/src/main/java/com/inteligr8/alfresco/asie/service/AbstractActionService.java +2 −10 Original line number Diff line number Diff line Loading @@ -12,17 +12,12 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.alfresco.model.ContentModel; import org.alfresco.repo.index.shard.Floc; import org.alfresco.repo.index.shard.Shard; import org.alfresco.repo.index.shard.ShardInstance; import org.alfresco.repo.index.shard.ShardRegistry; import org.alfresco.repo.index.shard.ShardState; import org.alfresco.service.cmr.repository.StoreRef; import org.alfresco.service.cmr.search.SearchParameters; import org.alfresco.service.cmr.search.SearchService; import org.alfresco.service.namespace.NamespaceService; import org.alfresco.service.namespace.QName; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; Loading @@ -45,9 +40,6 @@ public abstract class AbstractActionService { private final Logger logger = LoggerFactory.getLogger(this.getClass()); @Autowired private NamespaceService namespaceService; @Autowired private ApiService apiService; Loading @@ -58,10 +50,10 @@ public abstract class AbstractActionService { @Qualifier(Constants.QUALIFIER_ASIE) private ShardRegistry shardRegistry; @Value("${inteligr8.asie.default.concurrentQueueSize:64}") @Value("${inteligr8.asie.default.concurrentQueueSize}") private int concurrentQueueSize; @Value("${inteligr8.asie.default.concurrency:16}") @Value("${inteligr8.asie.default.concurrency}") private int concurrency; protected int getConcurrency() { Loading shared/src/main/java/com/inteligr8/alfresco/asie/service/AbstractNodeActionService.java +22 −8 Original line number Diff line number Diff line Loading @@ -8,6 +8,7 @@ import java.util.Map.Entry; import java.util.Set; import java.util.concurrent.Callable; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; Loading @@ -25,6 +26,7 @@ import org.alfresco.service.namespace.NamespaceService; import org.alfresco.service.namespace.QName; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Value; Loading @@ -41,7 +43,7 @@ import com.inteligr8.solr.model.Action; import com.inteligr8.solr.model.ActionResponse; import com.inteligr8.solr.model.BaseResponse; public abstract class AbstractNodeActionService { public abstract class AbstractNodeActionService implements DisposableBean { private final Logger logger = LoggerFactory.getLogger(this.getClass()); Loading @@ -61,10 +63,10 @@ public abstract class AbstractNodeActionService { @Qualifier(Constants.QUALIFIER_ASIE) private ShardRegistry shardRegistry; @Value("${inteligr8.asie.default.concurrentQueueSize:64}") @Value("${inteligr8.asie.default.concurrentQueueSize}") private int concurrentQueueSize; @Value("${inteligr8.asie.default.concurrency:16}") @Value("${inteligr8.asie.default.concurrency}") private int concurrency; protected int getConcurrency() { Loading @@ -79,6 +81,22 @@ public abstract class AbstractNodeActionService { protected abstract String getActionName(); @Override public void destroy() { ExecutorService executor = this.executorManager.get(this.getThreadNamePrefix()); if (executor != null) { this.logger.info("Shutting down throttled thread pool executor: {}", this.getThreadNamePrefix()); executor.shutdown(); } } private ThrottledThreadPoolExecutor getExecutor() { return this.executorManager.createThrottled( this.getThreadNamePrefix(), this.getConcurrency(), this.getConcurrency(), this.getConcurrentQueueSize(), 1L, TimeUnit.MINUTES); } /** * This method executes an action on the specified node in Solr using its * ACS unique database identifier. The callback handles all the return Loading Loading @@ -138,11 +156,7 @@ public abstract class AbstractNodeActionService { this.logger.debug("Will attempt to {} ACS node against {} shard instances: {}", this.getActionName(), eligibleInstances.size(), nodeDbId); CompositeFuture<Void> future = new CompositeFuture<>(); ThrottledThreadPoolExecutor executor = this.executorManager.createThrottled( this.getThreadNamePrefix(), this.getConcurrency(), this.getConcurrency(), this.getConcurrentQueueSize(), 1L, TimeUnit.MINUTES); ThrottledThreadPoolExecutor executor = this.getExecutor(); for (final com.inteligr8.alfresco.asie.model.ShardInstance instance : eligibleInstances) { this.logger.trace("Will attempt to {} ACS node against shard instance: {}: {}", this.getActionName(), nodeDbId, instance); Loading shared/src/main/java/com/inteligr8/alfresco/asie/service/AcsReconcileService.java +67 −71 Original line number Diff line number Diff line package com.inteligr8.alfresco.asie.service; import java.util.HashMap; import java.util.HashSet; import java.util.Map; import java.util.Collections; import java.util.Set; import java.util.concurrent.Callable; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; Loading @@ -26,7 +25,6 @@ import org.apache.commons.collections4.SetUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; Loading @@ -39,7 +37,7 @@ import com.inteligr8.alfresco.asie.util.CompositeFuture; import com.inteligr8.alfresco.asie.util.ThrottledThreadPoolExecutor; @Component public class AcsReconcileService implements InitializingBean, DisposableBean { public class AcsReconcileService implements DisposableBean { private final Logger logger = LoggerFactory.getLogger(this.getClass()); private final Logger reconcileLogger = LoggerFactory.getLogger("inteligr8.asie.reconcile"); Loading @@ -62,29 +60,35 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { @Autowired private ReindexService reindexService; @Value("${inteligr8.asie.reconciliation.nodesChunkSize:250}") @Autowired private ExecutorManager executorManager; @Value("${inteligr8.asie.reconciliation.nodesChunkSize}") private int nodesChunkSize; @Value("${inteligr8.asie.reconciliation.nodeTimeoutSeconds:10}") @Value("${inteligr8.asie.reconciliation.nodeTimeoutSeconds}") private int nodeTimeoutSeconds; @Value("${inteligr8.asie.reconciliation.concurrentQueueSize:64}") @Value("${inteligr8.asie.reconciliation.concurrentQueueSize}") private int concurrentQueueSize; @Value("${inteligr8.asie.reconciliation.concurrency:2}") @Value("${inteligr8.asie.reconciliation.concurrency}") private int concurrency; private ThrottledThreadPoolExecutor executor; @Override public void afterPropertiesSet() { this.executor = new ThrottledThreadPoolExecutor(this.concurrency, this.concurrency, this.concurrentQueueSize, 1L, TimeUnit.MINUTES, "solr-reconcile"); this.executor.prestartAllCoreThreads(); public void destroy() { ExecutorService executor = this.executorManager.get("solr-reconcile"); if (executor != null) { this.logger.info("Shutting down throttled thread pool executor: {}", "solr-reconcile"); executor.shutdown(); } } @Override public void destroy() { this.executor.shutdown(); private ThrottledThreadPoolExecutor getExecutor() { return this.executorManager.createThrottled( "solr-reconcile", this.concurrency, this.concurrency, this.concurrentQueueSize, 10L, TimeUnit.SECONDS); } /** Loading Loading @@ -205,6 +209,7 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { } CompositeFuture<Void> future = new CompositeFuture<>(); ThrottledThreadPoolExecutor executor = this.getExecutor(); for (long _nodeDbId = fromDbId; _nodeDbId < toDbId; _nodeDbId++) { final long nodeDbId = _nodeDbId; Loading @@ -222,7 +227,11 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { callback.reconciled(nodeDbId); if (reindexReconciled) reindex(nodeDbId, nodeRefs[dbIdIndex], callback, execTimeout, execUnit); reindex(nodeDbId, nodeRefs[dbIdIndex], callback); // purposefully forgetting about the returned future // its results will be logged // the reconcile thread will continue independently // the callback will lag return null; } }; Loading @@ -230,16 +239,16 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { callable = new Callable<Void>() { @Override public Void call() throws InterruptedException, TimeoutException { reconcile(nodeDbId, indexUnreconciled, callback, execTimeout, execUnit); reconcile(nodeDbId, indexUnreconciled, callback); return null; } }; } if (queueTimeout < 0L) { future.combine(this.executor.submit(callable, -1L, null)); future.combine(executor.submit(callable, -1L, null)); } else { future.combine(this.executor.submit(callable, queueTimeout, queueUnit)); future.combine(executor.submit(callable, queueTimeout, queueUnit)); } } Loading @@ -248,8 +257,7 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { public void reconcile(long nodeDbId, boolean index, ReconcileCallback callback, long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException { ReconcileCallback callback) throws InterruptedException, TimeoutException { NodeRef nodeRef = this.nodeService.getNodeRef(nodeDbId); if (nodeRef == null) { this.logger.trace("No such ACS node: {}; skipping ...", nodeDbId); Loading @@ -273,93 +281,81 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { this.reconcileLogger.info("UNRECONCILED: {} <=> {}", nodeDbId, nodeRef); callback.unreconciled(nodeDbId); } else { logger.debug("A node in the DB is not indexed in Solr; attempt to index: {}: {}", nodeDbId, nodeRef); this.index(nodeDbId, nodeRef, callback, execTimeout, execUnit); this.logger.debug("A node in the DB is not indexed in Solr; attempt to index: {}: {}", nodeDbId, nodeRef); this.index(nodeDbId, nodeRef, callback); // purposefully forgetting about the returned future // its results will be logged // the reconcile thread will continue independently // the callback will lag } } public void index(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback, long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException { Set<ShardInstance> syncHosts = new HashSet<>(); Set<ShardInstance> asyncHosts = new HashSet<>(); Map<ShardInstance, String> errorHosts = new HashMap<>(); public Future<Void> index(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback) throws InterruptedException, TimeoutException { IndexCallback indexCallback = new IndexCallback() { @Override public void success(ShardInstance instance) { reconcileLogger.info("INDEXED: {} <=> {} in {}", nodeDbId, nodeRef, instance); syncHosts.add(instance); if (callback != null) callback.processed(nodeDbId, Collections.singleton(instance), Collections.emptySet(), Collections.emptyMap()); } @Override public void scheduled(ShardInstance instance) { reconcileLogger.info("INDEXING: {} <=> {} in {}", nodeDbId, nodeRef, instance); asyncHosts.add(instance); if (callback != null) callback.processed(nodeDbId, Collections.emptySet(), Collections.singleton(instance), Collections.emptyMap()); } @Override public void error(ShardInstance instance, String message) { reconcileLogger.info("FAILED INDEX: {} <=> {} in {}", nodeDbId, nodeRef, instance); errorHosts.put(instance, message); if (callback != null) callback.processed(nodeDbId, Collections.emptySet(), Collections.emptySet(), Collections.singletonMap(instance, message)); } }; try { if (execTimeout < 0L) { this.indexService.index(nodeDbId, indexCallback).get(); } else { this.indexService.index(nodeDbId, indexCallback).get(execTimeout, execUnit); } } catch (ExecutionException ee) { throw new RuntimeException("An unexpected exception occurred: " + ee.getMessage(), ee); return this.indexService.index(nodeDbId, indexCallback); } if (callback != null) callback.processed(nodeDbId, syncHosts, asyncHosts, errorHosts); } public void reindex(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback, long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException { Set<ShardInstance> syncHosts = new HashSet<>(); Set<ShardInstance> asyncHosts = new HashSet<>(); Map<ShardInstance, String> errorHosts = new HashMap<>(); public Future<Void> reindex(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback) throws InterruptedException { ReindexCallback reindexCallback = new ReindexCallback() { @Override public void success(ShardInstance instance) { reconcileLogger.info("REINDEXED: {} <=> {} in {}", nodeDbId, nodeRef, instance); syncHosts.add(instance); if (callback != null) callback.processed(nodeDbId, Collections.singleton(instance), Collections.emptySet(), Collections.emptyMap()); } @Override public void scheduled(ShardInstance instance) { reconcileLogger.info("REINDEXING: {} <=> {} in {}", nodeDbId, nodeRef, instance); asyncHosts.add(instance); if (callback != null) callback.processed(nodeDbId, Collections.emptySet(), Collections.singleton(instance), Collections.emptyMap()); } @Override public void error(ShardInstance instance, String message) { reconcileLogger.info("FAILED REINDEX: {} <=> {} in {}", nodeDbId, nodeRef, instance); errorHosts.put(instance, message); if (callback != null) callback.processed(nodeDbId, Collections.emptySet(), Collections.emptySet(), Collections.singletonMap(instance, message)); } }; try { if (execTimeout < 0L) { this.reindexService.reindex(nodeDbId, reindexCallback).get(); } else { this.reindexService.reindex(nodeDbId, reindexCallback).get(execTimeout, execUnit); } } catch (ExecutionException ee) { throw new RuntimeException("An unexpected exception occurred: " + ee.getMessage(), ee); } if (callback != null) callback.processed(nodeDbId, syncHosts, asyncHosts, errorHosts); return this.reindexService.reindex(nodeDbId, reindexCallback); } private String formatForFts(QName qname) { Loading shared/src/main/java/com/inteligr8/alfresco/asie/service/ApiService.java +1 −1 Original line number Diff line number Diff line Loading @@ -43,7 +43,7 @@ public class ApiService implements InitializingBean { @Value("${inteligr8.asie.basePath}") private String solrBaseUrl; @Value("${inteligr8.asie.reconciliation.nodesChunkSize:250}") @Value("${inteligr8.asie.reconciliation.nodesChunkSize}") private int nodesChunkSize; @Override Loading Loading
shared/src/main/java/com/inteligr8/alfresco/asie/rest/AbstractActionWebScript.java +0 −6 Original line number Diff line number Diff line Loading @@ -9,14 +9,8 @@ import java.util.Map; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.alfresco.model.ContentModel; import org.alfresco.service.cmr.repository.InvalidNodeRefException; import org.alfresco.service.cmr.repository.NodeRef; import org.alfresco.service.cmr.repository.NodeService; import org.alfresco.service.cmr.repository.StoreRef; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.extensions.webscripts.WebScriptException; import org.springframework.extensions.webscripts.WebScriptRequest; import org.springframework.extensions.webscripts.WebScriptResponse; Loading
shared/src/main/java/com/inteligr8/alfresco/asie/service/AbstractActionService.java +2 −10 Original line number Diff line number Diff line Loading @@ -12,17 +12,12 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.alfresco.model.ContentModel; import org.alfresco.repo.index.shard.Floc; import org.alfresco.repo.index.shard.Shard; import org.alfresco.repo.index.shard.ShardInstance; import org.alfresco.repo.index.shard.ShardRegistry; import org.alfresco.repo.index.shard.ShardState; import org.alfresco.service.cmr.repository.StoreRef; import org.alfresco.service.cmr.search.SearchParameters; import org.alfresco.service.cmr.search.SearchService; import org.alfresco.service.namespace.NamespaceService; import org.alfresco.service.namespace.QName; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; Loading @@ -45,9 +40,6 @@ public abstract class AbstractActionService { private final Logger logger = LoggerFactory.getLogger(this.getClass()); @Autowired private NamespaceService namespaceService; @Autowired private ApiService apiService; Loading @@ -58,10 +50,10 @@ public abstract class AbstractActionService { @Qualifier(Constants.QUALIFIER_ASIE) private ShardRegistry shardRegistry; @Value("${inteligr8.asie.default.concurrentQueueSize:64}") @Value("${inteligr8.asie.default.concurrentQueueSize}") private int concurrentQueueSize; @Value("${inteligr8.asie.default.concurrency:16}") @Value("${inteligr8.asie.default.concurrency}") private int concurrency; protected int getConcurrency() { Loading
shared/src/main/java/com/inteligr8/alfresco/asie/service/AbstractNodeActionService.java +22 −8 Original line number Diff line number Diff line Loading @@ -8,6 +8,7 @@ import java.util.Map.Entry; import java.util.Set; import java.util.concurrent.Callable; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; Loading @@ -25,6 +26,7 @@ import org.alfresco.service.namespace.NamespaceService; import org.alfresco.service.namespace.QName; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Value; Loading @@ -41,7 +43,7 @@ import com.inteligr8.solr.model.Action; import com.inteligr8.solr.model.ActionResponse; import com.inteligr8.solr.model.BaseResponse; public abstract class AbstractNodeActionService { public abstract class AbstractNodeActionService implements DisposableBean { private final Logger logger = LoggerFactory.getLogger(this.getClass()); Loading @@ -61,10 +63,10 @@ public abstract class AbstractNodeActionService { @Qualifier(Constants.QUALIFIER_ASIE) private ShardRegistry shardRegistry; @Value("${inteligr8.asie.default.concurrentQueueSize:64}") @Value("${inteligr8.asie.default.concurrentQueueSize}") private int concurrentQueueSize; @Value("${inteligr8.asie.default.concurrency:16}") @Value("${inteligr8.asie.default.concurrency}") private int concurrency; protected int getConcurrency() { Loading @@ -79,6 +81,22 @@ public abstract class AbstractNodeActionService { protected abstract String getActionName(); @Override public void destroy() { ExecutorService executor = this.executorManager.get(this.getThreadNamePrefix()); if (executor != null) { this.logger.info("Shutting down throttled thread pool executor: {}", this.getThreadNamePrefix()); executor.shutdown(); } } private ThrottledThreadPoolExecutor getExecutor() { return this.executorManager.createThrottled( this.getThreadNamePrefix(), this.getConcurrency(), this.getConcurrency(), this.getConcurrentQueueSize(), 1L, TimeUnit.MINUTES); } /** * This method executes an action on the specified node in Solr using its * ACS unique database identifier. The callback handles all the return Loading Loading @@ -138,11 +156,7 @@ public abstract class AbstractNodeActionService { this.logger.debug("Will attempt to {} ACS node against {} shard instances: {}", this.getActionName(), eligibleInstances.size(), nodeDbId); CompositeFuture<Void> future = new CompositeFuture<>(); ThrottledThreadPoolExecutor executor = this.executorManager.createThrottled( this.getThreadNamePrefix(), this.getConcurrency(), this.getConcurrency(), this.getConcurrentQueueSize(), 1L, TimeUnit.MINUTES); ThrottledThreadPoolExecutor executor = this.getExecutor(); for (final com.inteligr8.alfresco.asie.model.ShardInstance instance : eligibleInstances) { this.logger.trace("Will attempt to {} ACS node against shard instance: {}: {}", this.getActionName(), nodeDbId, instance); Loading
shared/src/main/java/com/inteligr8/alfresco/asie/service/AcsReconcileService.java +67 −71 Original line number Diff line number Diff line package com.inteligr8.alfresco.asie.service; import java.util.HashMap; import java.util.HashSet; import java.util.Map; import java.util.Collections; import java.util.Set; import java.util.concurrent.Callable; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; Loading @@ -26,7 +25,6 @@ import org.apache.commons.collections4.SetUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; Loading @@ -39,7 +37,7 @@ import com.inteligr8.alfresco.asie.util.CompositeFuture; import com.inteligr8.alfresco.asie.util.ThrottledThreadPoolExecutor; @Component public class AcsReconcileService implements InitializingBean, DisposableBean { public class AcsReconcileService implements DisposableBean { private final Logger logger = LoggerFactory.getLogger(this.getClass()); private final Logger reconcileLogger = LoggerFactory.getLogger("inteligr8.asie.reconcile"); Loading @@ -62,29 +60,35 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { @Autowired private ReindexService reindexService; @Value("${inteligr8.asie.reconciliation.nodesChunkSize:250}") @Autowired private ExecutorManager executorManager; @Value("${inteligr8.asie.reconciliation.nodesChunkSize}") private int nodesChunkSize; @Value("${inteligr8.asie.reconciliation.nodeTimeoutSeconds:10}") @Value("${inteligr8.asie.reconciliation.nodeTimeoutSeconds}") private int nodeTimeoutSeconds; @Value("${inteligr8.asie.reconciliation.concurrentQueueSize:64}") @Value("${inteligr8.asie.reconciliation.concurrentQueueSize}") private int concurrentQueueSize; @Value("${inteligr8.asie.reconciliation.concurrency:2}") @Value("${inteligr8.asie.reconciliation.concurrency}") private int concurrency; private ThrottledThreadPoolExecutor executor; @Override public void afterPropertiesSet() { this.executor = new ThrottledThreadPoolExecutor(this.concurrency, this.concurrency, this.concurrentQueueSize, 1L, TimeUnit.MINUTES, "solr-reconcile"); this.executor.prestartAllCoreThreads(); public void destroy() { ExecutorService executor = this.executorManager.get("solr-reconcile"); if (executor != null) { this.logger.info("Shutting down throttled thread pool executor: {}", "solr-reconcile"); executor.shutdown(); } } @Override public void destroy() { this.executor.shutdown(); private ThrottledThreadPoolExecutor getExecutor() { return this.executorManager.createThrottled( "solr-reconcile", this.concurrency, this.concurrency, this.concurrentQueueSize, 10L, TimeUnit.SECONDS); } /** Loading Loading @@ -205,6 +209,7 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { } CompositeFuture<Void> future = new CompositeFuture<>(); ThrottledThreadPoolExecutor executor = this.getExecutor(); for (long _nodeDbId = fromDbId; _nodeDbId < toDbId; _nodeDbId++) { final long nodeDbId = _nodeDbId; Loading @@ -222,7 +227,11 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { callback.reconciled(nodeDbId); if (reindexReconciled) reindex(nodeDbId, nodeRefs[dbIdIndex], callback, execTimeout, execUnit); reindex(nodeDbId, nodeRefs[dbIdIndex], callback); // purposefully forgetting about the returned future // its results will be logged // the reconcile thread will continue independently // the callback will lag return null; } }; Loading @@ -230,16 +239,16 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { callable = new Callable<Void>() { @Override public Void call() throws InterruptedException, TimeoutException { reconcile(nodeDbId, indexUnreconciled, callback, execTimeout, execUnit); reconcile(nodeDbId, indexUnreconciled, callback); return null; } }; } if (queueTimeout < 0L) { future.combine(this.executor.submit(callable, -1L, null)); future.combine(executor.submit(callable, -1L, null)); } else { future.combine(this.executor.submit(callable, queueTimeout, queueUnit)); future.combine(executor.submit(callable, queueTimeout, queueUnit)); } } Loading @@ -248,8 +257,7 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { public void reconcile(long nodeDbId, boolean index, ReconcileCallback callback, long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException { ReconcileCallback callback) throws InterruptedException, TimeoutException { NodeRef nodeRef = this.nodeService.getNodeRef(nodeDbId); if (nodeRef == null) { this.logger.trace("No such ACS node: {}; skipping ...", nodeDbId); Loading @@ -273,93 +281,81 @@ public class AcsReconcileService implements InitializingBean, DisposableBean { this.reconcileLogger.info("UNRECONCILED: {} <=> {}", nodeDbId, nodeRef); callback.unreconciled(nodeDbId); } else { logger.debug("A node in the DB is not indexed in Solr; attempt to index: {}: {}", nodeDbId, nodeRef); this.index(nodeDbId, nodeRef, callback, execTimeout, execUnit); this.logger.debug("A node in the DB is not indexed in Solr; attempt to index: {}: {}", nodeDbId, nodeRef); this.index(nodeDbId, nodeRef, callback); // purposefully forgetting about the returned future // its results will be logged // the reconcile thread will continue independently // the callback will lag } } public void index(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback, long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException { Set<ShardInstance> syncHosts = new HashSet<>(); Set<ShardInstance> asyncHosts = new HashSet<>(); Map<ShardInstance, String> errorHosts = new HashMap<>(); public Future<Void> index(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback) throws InterruptedException, TimeoutException { IndexCallback indexCallback = new IndexCallback() { @Override public void success(ShardInstance instance) { reconcileLogger.info("INDEXED: {} <=> {} in {}", nodeDbId, nodeRef, instance); syncHosts.add(instance); if (callback != null) callback.processed(nodeDbId, Collections.singleton(instance), Collections.emptySet(), Collections.emptyMap()); } @Override public void scheduled(ShardInstance instance) { reconcileLogger.info("INDEXING: {} <=> {} in {}", nodeDbId, nodeRef, instance); asyncHosts.add(instance); if (callback != null) callback.processed(nodeDbId, Collections.emptySet(), Collections.singleton(instance), Collections.emptyMap()); } @Override public void error(ShardInstance instance, String message) { reconcileLogger.info("FAILED INDEX: {} <=> {} in {}", nodeDbId, nodeRef, instance); errorHosts.put(instance, message); if (callback != null) callback.processed(nodeDbId, Collections.emptySet(), Collections.emptySet(), Collections.singletonMap(instance, message)); } }; try { if (execTimeout < 0L) { this.indexService.index(nodeDbId, indexCallback).get(); } else { this.indexService.index(nodeDbId, indexCallback).get(execTimeout, execUnit); } } catch (ExecutionException ee) { throw new RuntimeException("An unexpected exception occurred: " + ee.getMessage(), ee); return this.indexService.index(nodeDbId, indexCallback); } if (callback != null) callback.processed(nodeDbId, syncHosts, asyncHosts, errorHosts); } public void reindex(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback, long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException { Set<ShardInstance> syncHosts = new HashSet<>(); Set<ShardInstance> asyncHosts = new HashSet<>(); Map<ShardInstance, String> errorHosts = new HashMap<>(); public Future<Void> reindex(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback) throws InterruptedException { ReindexCallback reindexCallback = new ReindexCallback() { @Override public void success(ShardInstance instance) { reconcileLogger.info("REINDEXED: {} <=> {} in {}", nodeDbId, nodeRef, instance); syncHosts.add(instance); if (callback != null) callback.processed(nodeDbId, Collections.singleton(instance), Collections.emptySet(), Collections.emptyMap()); } @Override public void scheduled(ShardInstance instance) { reconcileLogger.info("REINDEXING: {} <=> {} in {}", nodeDbId, nodeRef, instance); asyncHosts.add(instance); if (callback != null) callback.processed(nodeDbId, Collections.emptySet(), Collections.singleton(instance), Collections.emptyMap()); } @Override public void error(ShardInstance instance, String message) { reconcileLogger.info("FAILED REINDEX: {} <=> {} in {}", nodeDbId, nodeRef, instance); errorHosts.put(instance, message); if (callback != null) callback.processed(nodeDbId, Collections.emptySet(), Collections.emptySet(), Collections.singletonMap(instance, message)); } }; try { if (execTimeout < 0L) { this.reindexService.reindex(nodeDbId, reindexCallback).get(); } else { this.reindexService.reindex(nodeDbId, reindexCallback).get(execTimeout, execUnit); } } catch (ExecutionException ee) { throw new RuntimeException("An unexpected exception occurred: " + ee.getMessage(), ee); } if (callback != null) callback.processed(nodeDbId, syncHosts, asyncHosts, errorHosts); return this.reindexService.reindex(nodeDbId, reindexCallback); } private String formatForFts(QName qname) { Loading
shared/src/main/java/com/inteligr8/alfresco/asie/service/ApiService.java +1 −1 Original line number Diff line number Diff line Loading @@ -43,7 +43,7 @@ public class ApiService implements InitializingBean { @Value("${inteligr8.asie.basePath}") private String solrBaseUrl; @Value("${inteligr8.asie.reconciliation.nodesChunkSize:250}") @Value("${inteligr8.asie.reconciliation.nodesChunkSize}") private int nodesChunkSize; @Override Loading