Updates for allowing commands to be cancelled, within single processes is complete. Now need to add a spray to the cluster when cancel is called as well as forcing indexreaders to exit. Exception handling may need work.
Project: http://git-wip-us.apache.org/repos/asf/incubator-blur/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-blur/commit/1bbb14b2 Tree: http://git-wip-us.apache.org/repos/asf/incubator-blur/tree/1bbb14b2 Diff: http://git-wip-us.apache.org/repos/asf/incubator-blur/diff/1bbb14b2 Branch: refs/heads/master Commit: 1bbb14b2d76d233a99774a8a20df702b7b7d94de Parents: 0238b4a Author: Aaron McCurry <[email protected]> Authored: Fri Oct 31 10:37:07 2014 -0400 Committer: Aaron McCurry <[email protected]> Committed: Fri Oct 31 10:37:07 2014 -0400 ---------------------------------------------------------------------- .../apache/blur/command/ArgumentOverlay.java | 20 ++++ .../apache/blur/command/BaseCommandManager.java | 109 ++++++++++++------- .../java/org/apache/blur/command/Command.java | 11 ++ .../org/apache/blur/command/CommandRunner.java | 4 +- .../blur/command/ControllerClusterContext.java | 8 +- .../blur/command/ControllerCommandManager.java | 2 +- .../org/apache/blur/command/ResponseFuture.java | 16 ++- .../blur/command/ShardCommandManager.java | 17 +-- .../apache/blur/command/TimeoutException.java | 10 +- .../apache/blur/server/FilteredBlurServer.java | 22 ++-- .../blur/thrift/BlurControllerServer.java | 15 ++- .../org/apache/blur/thrift/BlurShardServer.java | 18 +-- .../blur/command/ShardCommandManagerTest.java | 81 +++++++++++--- 13 files changed, 225 insertions(+), 108 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java b/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java index e051b93..493f1f8 100644 --- a/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java +++ b/blur-core/src/main/java/org/apache/blur/command/ArgumentOverlay.java @@ -18,12 +18,17 @@ package org.apache.blur.command; import java.lang.reflect.Field; import java.util.Map; +import java.util.UUID; import org.apache.blur.command.annotation.OptionalArgument; import org.apache.blur.command.annotation.RequiredArgument; +import org.apache.blur.log.Log; +import org.apache.blur.log.LogFactory; public class ArgumentOverlay { + private static final Log LOG = LogFactory.getLog(ArgumentOverlay.class); + private final Map<String, ? extends Object> _args; public ArgumentOverlay(BlurObject args, BlurObjectSerDe serDe) { @@ -33,14 +38,29 @@ public class ArgumentOverlay { public <T> Command<T> setup(Command<T> command) { Class<?> clazz = command.getClass(); setupInternal(clazz, command); + return setupCommandExecutionId(command); + } + + private <T> Command<T> setupCommandExecutionId(Command<T> command) { + String commandExecutionId = command.getCommandExecutionId(); + if (commandExecutionId == null) { + commandExecutionId = UUID.randomUUID().toString(); + LOG.info("Command execution id [{0}] has been assigned to [{1}]", commandExecutionId, command); + command.setCommandExecutionId(commandExecutionId); + } return command; } private void setupInternal(Class<?> clazz, Command<?> command) { if (clazz.equals(Command.class)) { + mapValuesToFields(clazz, command); return; } setupInternal(clazz.getSuperclass(), command); + mapValuesToFields(clazz, command); + } + + private void mapValuesToFields(Class<?> clazz, Command<?> command) { Field[] declaredFields = clazz.getDeclaredFields(); for (Field field : declaredFields) { RequiredArgument requiredArgument = field.getAnnotation(RequiredArgument.class); http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java b/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java index ba6d4be..6ccded6 100644 --- a/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java +++ b/blur-core/src/main/java/org/apache/blur/command/BaseCommandManager.java @@ -16,6 +16,7 @@ import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.Properties; +import java.util.Random; import java.util.Set; import java.util.Timer; import java.util.TimerTask; @@ -66,13 +67,15 @@ public abstract class BaseCommandManager implements Closeable { private static final String META_INF_SERVICES_ORG_APACHE_BLUR_COMMAND_COMMANDS = "META-INF/services/org.apache.blur.command.Commands"; private static final Log LOG = LogFactory.getLog(BaseCommandManager.class); - private final ExecutorService _executorService; + private final ExecutorService _executorServiceWorker; private final ExecutorService _executorServiceDriver; + private final Random _random = new Random(); protected final Map<String, BigInteger> _commandLoadTime = new ConcurrentHashMap<String, BigInteger>(); protected final Map<String, Command<?>> _command = new ConcurrentHashMap<String, Command<?>>(); protected final Map<Class<? extends Command<?>>, String> _commandNameLookup = new ConcurrentHashMap<Class<? extends Command<?>>, String>(); - protected final ConcurrentMap<ExecutionId, ResponseFuture> _runningMap = new MapMaker().makeMap(); + protected final ConcurrentMap<Long, ResponseFuture<?>> _driverRunningMap = new MapMaker().makeMap(); + protected final ConcurrentMap<Long, ResponseFuture<?>> _workerRunningMap = new MapMaker().makeMap(); protected final long _connectionTimeout; protected final File _tmpPath; protected final String _commandPath; @@ -89,11 +92,12 @@ public abstract class BaseCommandManager implements Closeable { lookForCommandsToRegisterInClassPath(); _tmpPath = tmpPath; _commandPath = commandPath; - _executorService = Executors.newThreadPool("command-worker-", workerThreadCount); + _executorServiceWorker = Executors.newThreadPool("command-worker-", workerThreadCount); _executorServiceDriver = Executors.newThreadPool("command-driver-", driverThreadCount); _connectionTimeout = connectionTimeout / 2; _timer = new Timer("BaseCommandManager-Timer", true); - _timer.schedule(getTimerTaskForRemovalOfOldCommands(), _pollingPeriod, _pollingPeriod); + _timer.schedule(getTimerTaskForRemovalOfOldCommands(_driverRunningMap), _pollingPeriod, _pollingPeriod); + _timer.schedule(getTimerTaskForRemovalOfOldCommands(_workerRunningMap), _pollingPeriod, _pollingPeriod); if (_tmpPath == null || _commandPath == null) { LOG.info("Tmp Path [{0}] or Command Path [{1}] is null so the automatic command reload will be disabled.", _tmpPath, _commandPath); @@ -103,30 +107,45 @@ public abstract class BaseCommandManager implements Closeable { } } - public void cancel(ExecutionId executionId) { - ResponseFuture responseFuture = _runningMap.get(executionId); - responseFuture.cancel(true); + public void cancelCommand(String commandExecutionId) { + LOG.info("Trying to cancel command [{0}]", commandExecutionId); + cancelAllExecuting(commandExecutionId, _workerRunningMap); + cancelAllExecuting(commandExecutionId, _driverRunningMap); } - public List<String> commandStatusList(CommandStatusStateEnum commandStatus) { - throw new RuntimeException("Not implemented."); + private void cancelAllExecuting(String commandExecutionId, ConcurrentMap<Long, ResponseFuture<?>> runningMap) { + for (Entry<Long, ResponseFuture<?>> e : runningMap.entrySet()) { + Long instanceExecutionId = e.getKey(); + ResponseFuture<?> future = e.getValue(); + Command<?> commandExecuting = future.getCommandExecuting(); + if (commandExecuting.getCommandExecutionId().equals(commandExecutionId)) { + LOG.info("Canceling Command with executing id [{0}] command [{1}]", instanceExecutionId, commandExecuting); + future.cancel(true); + } + } } - public void cancel(String executionId) { - cancel(new ExecutionId(executionId)); + public List<String> commandStatusList(CommandStatusStateEnum commandStatus) { + throw new RuntimeException("Not implemented."); } - private TimerTask getTimerTaskForRemovalOfOldCommands() { + private TimerTask getTimerTaskForRemovalOfOldCommands(final Map<Long, ResponseFuture<?>> runningMap) { return new TimerTask() { @Override public void run() { - Set<Entry<ExecutionId, ResponseFuture>> entrySet = _runningMap.entrySet(); - for (Entry<ExecutionId, ResponseFuture> e : entrySet) { - ExecutionId executionId = e.getKey(); - ResponseFuture responseFuture = e.getValue(); + Set<Entry<Long, ResponseFuture<?>>> entrySet = runningMap.entrySet(); + for (Entry<Long, ResponseFuture<?>> e : entrySet) { + Long instanceExecutionId = e.getKey(); + ResponseFuture<?> responseFuture = e.getValue(); if (!responseFuture.isRunning() && responseFuture.hasExpired()) { - LOG.info("Removing old execution id [{0}]", executionId); - _runningMap.remove(executionId); + Command<?> commandExecuting = responseFuture.getCommandExecuting(); + String commandExecutionId = null; + if (commandExecuting != null) { + commandExecutionId = commandExecuting.getCommandExecutionId(); + } + LOG.info("Removing old execution instance id [{0}] with command execution id of [{1}]", + instanceExecutionId, commandExecutionId); + runningMap.remove(instanceExecutionId); } } } @@ -284,10 +303,11 @@ public abstract class BaseCommandManager implements Closeable { } } - public Response reconnect(ExecutionId executionId) throws IOException, TimeoutException { - Future<Response> future = _runningMap.get(executionId); + @SuppressWarnings("unchecked") + public Response reconnect(Long instanceExecutionId) throws IOException, TimeoutException { + Future<Response> future = (Future<Response>) _driverRunningMap.get(instanceExecutionId); if (future == null) { - throw new IOException("Command id [" + executionId + "] did not find any executing commands."); + throw new IOException("Execution instance id [" + instanceExecutionId + "] did not find any executing commands."); } try { return future.get(_connectionTimeout, TimeUnit.MILLISECONDS); @@ -298,18 +318,17 @@ public abstract class BaseCommandManager implements Closeable { } catch (ExecutionException e) { throw new IOException(e.getCause()); } catch (java.util.concurrent.TimeoutException e) { - LOG.info("Timeout of command [{0}]", executionId); - throw new TimeoutException(executionId); + LOG.info("Timeout of command [{0}]", instanceExecutionId); + throw new TimeoutException(instanceExecutionId); } } - protected Response submitDriverCallable(Callable<Response> callable) throws IOException, TimeoutException, - ExceptionCollector { - ExecutionContext executionContext = ExecutionContext.create(); - Future<Response> future = _executorServiceDriver.submit(executionContext.wrapCallable(callable)); - executionContext.registerDriverFuture(future); - ExecutionId executionId = executionContext.getExecutionId(); - _runningMap.put(executionId, new ResponseFuture(_runningCacheTombstoneTime, future)); + protected Response submitDriverCallable(Callable<Response> callable, Command<?> commandExecuting) throws IOException, + TimeoutException, ExceptionCollector { + Future<Response> future = _executorServiceDriver.submit(callable); + Long instanceExecutionId = getInstanceExecutionId(); + _driverRunningMap.put(instanceExecutionId, new ResponseFuture<Response>(_runningCacheTombstoneTime, future, + commandExecuting)); try { return future.get(_connectionTimeout, TimeUnit.MILLISECONDS); } catch (CancellationException e) { @@ -323,21 +342,37 @@ public abstract class BaseCommandManager implements Closeable { } throw new IOException(cause); } catch (java.util.concurrent.TimeoutException e) { - LOG.info("Timeout of command [{0}]", executionId); - throw new TimeoutException(executionId); + LOG.info("Timeout of command [{0}]", instanceExecutionId); + throw new TimeoutException(instanceExecutionId); + } + } + + private Long getInstanceExecutionId() { + synchronized (_random) { + while (true) { + Long id = _random.nextLong(); + if (_driverRunningMap.containsKey(id)) { + continue; + } + if (_workerRunningMap.containsKey(id)) { + continue; + } + return id; + } } } - protected <T> Future<T> submitToExecutorService(Callable<T> callable) { - ExecutionContext executionContext = ExecutionContext.get(); - Future<T> future = _executorService.submit(executionContext.wrapCallable(callable)); - executionContext.registerFuture(future); + protected <T> Future<T> submitToExecutorService(Callable<T> callable, Command<?> commandExecuting) { + Future<T> future = _executorServiceWorker.submit(callable); + Long instanceExecutionId = getInstanceExecutionId(); + _workerRunningMap.put(instanceExecutionId, new ResponseFuture<T>(_runningCacheTombstoneTime, future, + commandExecuting)); return future; } @Override public void close() throws IOException { - _executorService.shutdownNow(); + _executorServiceWorker.shutdownNow(); _executorServiceDriver.shutdownNow(); if (_timer != null) { _timer.cancel(); http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/Command.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/command/Command.java b/blur-core/src/main/java/org/apache/blur/command/Command.java index 89b6e0e..9cf2719 100644 --- a/blur-core/src/main/java/org/apache/blur/command/Command.java +++ b/blur-core/src/main/java/org/apache/blur/command/Command.java @@ -30,6 +30,9 @@ import org.apache.blur.thrift.generated.Blur.Iface; public abstract class Command<R> implements Cloneable { + @OptionalArgument("The ") + private String commandExecutionId; + public abstract String getName(); public abstract R run() throws IOException; @@ -117,4 +120,12 @@ public abstract class Command<R> implements Cloneable { return builder.toString(); } + public String getCommandExecutionId() { + return commandExecutionId; + } + + public void setCommandExecutionId(String commandExecutionId) { + this.commandExecutionId = commandExecutionId; + } + } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java b/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java index f88d6c2..38c38aa 100644 --- a/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java +++ b/blur-core/src/main/java/org/apache/blur/command/CommandRunner.java @@ -209,7 +209,7 @@ public class CommandRunner { ClientPool clientPool = BlurClientManager.getClientPool(); Client client = clientPool.getClient(connection); try { - String executionId = null; + Long executionId = null; Response response; INNER: while (true) { try { @@ -220,7 +220,7 @@ public class CommandRunner { } break INNER; } catch (TimeoutException te) { - executionId = te.getExecutionId(); + executionId = te.getInstanceExecutionId(); } } return CommandUtil.fromThriftResponseToObject(response); http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java b/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java index 29c4c70..13b1589 100644 --- a/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java +++ b/blur-core/src/main/java/org/apache/blur/command/ControllerClusterContext.java @@ -126,7 +126,7 @@ public class ControllerClusterContext extends ClusterContext implements Closeabl Map<Shard, Object> shardToValue = CommandUtil.fromThriftSupportedObjects(shardToThriftValue, _serDe); return (Map<Shard, T>) shardToValue; } - }); + }, command); for (Shard shard : getShardsOnServer(server, tables, shards)) { futureMap.put(shard, new ShardResultFuture<T>(shard, future)); } @@ -149,7 +149,7 @@ public class ControllerClusterContext extends ClusterContext implements Closeabl protected static Response waitForResponse(Client client, Command<?> command, Arguments arguments) throws TException { // TODO This should likely be changed to run of a AtomicBoolean used for // the status of commands. - String executionId = null; + Long executionId = null; while (true) { try { if (executionId == null) { @@ -160,7 +160,7 @@ public class ControllerClusterContext extends ClusterContext implements Closeabl } catch (BlurException e) { throw e; } catch (TimeoutException e) { - executionId = e.getExecutionId(); + executionId = e.getInstanceExecutionId(); LOG.info("Execution fetch timed out, reconnecting using [{0}].", executionId); } catch (TException e) { throw e; @@ -212,7 +212,7 @@ public class ControllerClusterContext extends ClusterContext implements Closeabl Object thriftObject = CommandUtil.toObject(valueObject); return (T) _serDe.fromSupportedThriftObject(thriftObject); } - }); + }, command); futureMap.put(server, future); } return futureMap; http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java b/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java index 31901f8..386bae8 100644 --- a/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java +++ b/blur-core/src/main/java/org/apache/blur/command/ControllerCommandManager.java @@ -77,7 +77,7 @@ public class ControllerCommandManager extends BaseCommandManager { throw new IOException("Command type of [" + command.getClass() + "] not supported."); } - }); + }, command); } private CombiningContext getCombiningContext(final TableContextFactory tableContextFactory) { http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java b/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java index ca75b3e..eb9b30d 100644 --- a/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java +++ b/blur-core/src/main/java/org/apache/blur/command/ResponseFuture.java @@ -22,15 +22,21 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicLong; -public class ResponseFuture implements Future<Response> { +public class ResponseFuture<T> implements Future<T> { - private final Future<Response> _future; + private final Future<T> _future; private final AtomicLong _timeWhenNotRunningObserved = new AtomicLong(); private final long _tombstone; + private final Command<?> _commandExecuting; - public ResponseFuture(long tombstone, Future<Response> future) { + public ResponseFuture(long tombstone, Future<T> future, Command<?> commandExecuting) { _tombstone = tombstone; _future = future; + _commandExecuting = commandExecuting; + } + + public Command<?> getCommandExecuting() { + return _commandExecuting; } public boolean cancel(boolean mayInterruptIfRunning) { @@ -45,11 +51,11 @@ public class ResponseFuture implements Future<Response> { return _future.isDone(); } - public Response get() throws InterruptedException, ExecutionException { + public T get() throws InterruptedException, ExecutionException { return _future.get(); } - public Response get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { + public T get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { return _future.get(timeout, unit); } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java b/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java index ec3b6ab..cbf218f 100644 --- a/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java +++ b/blur-core/src/main/java/org/apache/blur/command/ShardCommandManager.java @@ -51,10 +51,10 @@ public class ShardCommandManager extends BaseCommandManager { public Response execute(final TableContextFactory tableContextFactory, final String commandName, final ArgumentOverlay argumentOverlay) throws IOException, TimeoutException, ExceptionCollector { final ShardServerContext shardServerContext = getShardServerContext(); + final Command<?> command = getCommandObject(commandName, argumentOverlay); Callable<Response> callable = new Callable<Response>() { @Override public Response call() throws Exception { - Command<?> command = getCommandObject(commandName, argumentOverlay); if (command == null) { throw new IOException("Command with name [" + commandName + "] not found."); } @@ -80,7 +80,7 @@ public class ShardCommandManager extends BaseCommandManager { }; } }; - return submitDriverCallable(callable); + return submitDriverCallable(callable, command); } private ShardServerContext getShardServerContext() { @@ -137,16 +137,17 @@ public class ShardCommandManager extends BaseCommandManager { } final BlurIndex blurIndex = e.getValue(); Callable<Object> callable; - if (command instanceof IndexRead) { - final IndexRead<?> readCommand = (IndexRead<?>) command.clone(); + Command<?> clone = command.clone(); + if (clone instanceof IndexRead) { + final IndexRead<?> readCommand = (IndexRead<?>) clone; callable = getCallable(shardServerContext, tableContextFactory, table, shard, blurIndex, readCommand); - } else if (command instanceof ServerRead) { - final ServerRead<?, ?> readCombiningCommand = (ServerRead<?, ?>) command.clone(); + } else if (clone instanceof ServerRead) { + final ServerRead<?, ?> readCombiningCommand = (ServerRead<?, ?>) clone; callable = getCallable(shardServerContext, tableContextFactory, table, shard, blurIndex, readCombiningCommand); } else { - throw new IOException("Command type of [" + command.getClass() + "] not supported."); + throw new IOException("Command type of [" + clone.getClass() + "] not supported."); } - Future<Object> future = submitToExecutorService(callable); + Future<Object> future = submitToExecutorService(callable, clone); futureMap.put(shard, future); } } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java b/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java index a0b091f..9f85bb2 100644 --- a/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java +++ b/blur-core/src/main/java/org/apache/blur/command/TimeoutException.java @@ -19,14 +19,14 @@ package org.apache.blur.command; @SuppressWarnings("serial") public class TimeoutException extends Exception { - private final ExecutionId _executionId; + private final Long _executionInstanceId; - public TimeoutException(ExecutionId executionId) { - _executionId = executionId; + public TimeoutException(Long executionInstanceId) { + _executionInstanceId = executionInstanceId; } - public ExecutionId getExecutionId() { - return _executionId; + public Long getInstanceExecutionId() { + return _executionInstanceId; } } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java b/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java index 2154916..9318f95 100644 --- a/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java +++ b/blur-core/src/main/java/org/apache/blur/server/FilteredBlurServer.java @@ -256,11 +256,6 @@ public class FilteredBlurServer implements Iface { } @Override - public Response reconnect(String executionId) throws BlurException, TimeoutException, TException { - return _iface.reconnect(executionId); - } - - @Override public void refresh() throws TException { _iface.refresh(); } @@ -272,18 +267,23 @@ public class FilteredBlurServer implements Iface { } @Override - public CommandStatus commandStatus(String executionId) throws BlurException, TException { - return _iface.commandStatus(executionId); + public List<CommandDescriptor> listInstalledCommands() throws BlurException, TException { + return _iface.listInstalledCommands(); + } + + @Override + public Response reconnect(long instanceExecutionId) throws BlurException, TimeoutException, TException { + return _iface.reconnect(instanceExecutionId); } @Override - public void commandCancel(String executionId) throws BlurException, TException { - _iface.commandCancel(executionId); + public CommandStatus commandStatus(String commandExecutionId) throws BlurException, TException { + return _iface.commandStatus(commandExecutionId); } @Override - public List<CommandDescriptor> listInstalledCommands() throws BlurException, TException { - return _iface.listInstalledCommands(); + public void commandCancel(String commandExecutionId) throws BlurException, TException { + _iface.commandCancel(commandExecutionId); } } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java b/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java index 477e0d0..23f719c 100644 --- a/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java +++ b/blur-core/src/main/java/org/apache/blur/thrift/BlurControllerServer.java @@ -53,7 +53,6 @@ import org.apache.blur.command.BlurObject; import org.apache.blur.command.BlurObjectSerDe; import org.apache.blur.command.CommandUtil; import org.apache.blur.command.ControllerCommandManager; -import org.apache.blur.command.ExecutionId; import org.apache.blur.command.Response; import org.apache.blur.command.Server; import org.apache.blur.command.Shard; @@ -1525,7 +1524,7 @@ public class BlurControllerServer extends TableAdmin implements Iface { return CommandUtil.fromObjectToThrift(response, _serDe); } catch (Exception e) { if (e instanceof org.apache.blur.command.TimeoutException) { - throw new TimeoutException(((org.apache.blur.command.TimeoutException) e).getExecutionId().getId()); + throw new TimeoutException(((org.apache.blur.command.TimeoutException) e).getInstanceExecutionId()); } LOG.error("Unknown error while trying to execute command [{0}]", e, commandName); if (e instanceof BlurException) { @@ -1623,16 +1622,16 @@ public class BlurControllerServer extends TableAdmin implements Iface { } @Override - public org.apache.blur.thrift.generated.Response reconnect(String executionId) throws BlurException, + public org.apache.blur.thrift.generated.Response reconnect(long instanceExecutionId) throws BlurException, TimeoutException, TException { try { - Response response = _commandManager.reconnect(new ExecutionId(executionId)); + Response response = _commandManager.reconnect(instanceExecutionId); return CommandUtil.fromObjectToThrift(response, _serDe); } catch (Exception e) { if (e instanceof org.apache.blur.command.TimeoutException) { - throw new TimeoutException(((org.apache.blur.command.TimeoutException) e).getExecutionId().getId()); + throw new TimeoutException(((org.apache.blur.command.TimeoutException) e).getInstanceExecutionId()); } - LOG.error("Unknown error while trying to reconnect to executing command [{0}]", e, executionId); + LOG.error("Unknown error while trying to reconnect to executing command [{0}]", e, instanceExecutionId); if (e instanceof BlurException) { throw (BlurException) e; } @@ -1665,12 +1664,12 @@ public class BlurControllerServer extends TableAdmin implements Iface { } @Override - public CommandStatus commandStatus(String executionId) throws BlurException, TException { + public CommandStatus commandStatus(String commandExecutionId) throws BlurException, TException { throw new BException("Not Implemented"); } @Override - public void commandCancel(String executionId) throws BlurException, TException { + public void commandCancel(String commandExecutionId) throws BlurException, TException { throw new BException("Not Implemented"); } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java b/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java index c1a5b16..135feaf 100644 --- a/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java +++ b/blur-core/src/main/java/org/apache/blur/thrift/BlurShardServer.java @@ -36,7 +36,6 @@ import org.apache.blur.command.BlurObject; import org.apache.blur.command.BlurObjectSerDe; import org.apache.blur.command.CommandStatusStateEnum; import org.apache.blur.command.CommandUtil; -import org.apache.blur.command.ExecutionId; import org.apache.blur.command.Response; import org.apache.blur.command.ShardCommandManager; import org.apache.blur.concurrent.Executors; @@ -617,7 +616,7 @@ public class BlurShardServer extends TableAdmin implements Iface { return CommandUtil.fromObjectToThrift(response, _serDe); } catch (Exception e) { if (e instanceof org.apache.blur.command.TimeoutException) { - throw new TimeoutException(((org.apache.blur.command.TimeoutException) e).getExecutionId().getId()); + throw new TimeoutException(((org.apache.blur.command.TimeoutException) e).getInstanceExecutionId()); } LOG.error("Unknown error while trying to execute command [{0}] for table [{1}]", e, commandName); if (e instanceof BlurException) { @@ -632,16 +631,16 @@ public class BlurShardServer extends TableAdmin implements Iface { } @Override - public org.apache.blur.thrift.generated.Response reconnect(String executionId) throws BlurException, + public org.apache.blur.thrift.generated.Response reconnect(long instanceExecutionId) throws BlurException, TimeoutException, TException { try { - Response response = _commandManager.reconnect(new ExecutionId(executionId)); + Response response = _commandManager.reconnect(instanceExecutionId); return CommandUtil.fromObjectToThrift(response, _serDe); } catch (Exception e) { if (e instanceof org.apache.blur.command.TimeoutException) { - throw new TimeoutException(((org.apache.blur.command.TimeoutException) e).getExecutionId().getId()); + throw new TimeoutException(((org.apache.blur.command.TimeoutException) e).getInstanceExecutionId()); } - LOG.error("Unknown error while trying to reconnect to executing command [{0}]", e, executionId); + LOG.error("Unknown error while trying to reconnect to executing command [{0}]", e, instanceExecutionId); if (e instanceof BlurException) { throw (BlurException) e; } @@ -679,14 +678,14 @@ public class BlurShardServer extends TableAdmin implements Iface { } @Override - public CommandStatus commandStatus(String executionId) throws BlurException, TException { + public CommandStatus commandStatus(String commandExecutionId) throws BlurException, TException { throw new BException("Not Implemented"); } @Override - public void commandCancel(String executionId) throws BlurException, TException { + public void commandCancel(String commandExecutionId) throws BlurException, TException { try { - _commandManager.cancel(executionId); + _commandManager.cancelCommand(commandExecutionId); } catch (Exception e) { throw new BException(e.getMessage(), e); } @@ -695,4 +694,5 @@ public class BlurShardServer extends TableAdmin implements Iface { private CommandStatusStateEnum toCommandStatus(CommandStatusState state) { return CommandStatusStateEnum.valueOf(state.name()); } + } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/1bbb14b2/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java ---------------------------------------------------------------------- diff --git a/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java b/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java index 83dd232..d2b4eed 100644 --- a/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java +++ b/blur-core/src/test/java/org/apache/blur/command/ShardCommandManagerTest.java @@ -30,6 +30,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.SortedSet; +import java.util.concurrent.CancellationException; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -40,6 +41,8 @@ import org.apache.blur.server.IndexSearcherClosable; import org.apache.blur.server.ShardContext; import org.apache.blur.server.TableContext; import org.apache.blur.server.TableContextFactory; +import org.apache.blur.thrift.generated.Arguments; +import org.apache.blur.thrift.generated.BlurException; import org.apache.blur.thrift.generated.RowMutation; import org.apache.blur.thrift.generated.ShardState; import org.apache.blur.thrift.generated.TableDescriptor; @@ -183,7 +186,7 @@ public class ShardCommandManagerTest { @Test public void testShardCommandManagerNormalWait() throws IOException, TimeoutException, ExceptionCollector { Response response; - ExecutionId executionId = null; + Long instanceExecutionId = null; BlurObject args = new BlurObject(); args.put("table", "test"); @@ -197,15 +200,15 @@ public class ShardCommandManagerTest { fail(); } try { - if (executionId == null) { + if (instanceExecutionId == null) { TableContextFactory tableContextFactory = getTableContextFactory(); response = _manager.execute(tableContextFactory, "wait", argumentOverlay); } else { - response = _manager.reconnect(executionId); + response = _manager.reconnect(instanceExecutionId); } break; } catch (TimeoutException te) { - executionId = te.getExecutionId(); + instanceExecutionId = te.getInstanceExecutionId(); } } System.out.println(response); @@ -229,24 +232,66 @@ public class ShardCommandManagerTest { } @Test - public void testShardCommandManagerNormalWithCancel() throws IOException, TimeoutException, ExceptionCollector { - Response response; - ExecutionId executionId = null; + public void testShardCommandManagerNormalWithCancel() throws IOException, TimeoutException, ExceptionCollector, + BlurException, InterruptedException { - BlurObject args = new BlurObject(); - args.put("table", "test"); - args.put("seconds", 5); + String commandExecutionId = "TEST_COMMAND_ID1"; - ArgumentOverlay argumentOverlay = new ArgumentOverlay(args, new BlurObjectSerDe()); + BlurObjectSerDe serDe = new BlurObjectSerDe(); + WaitForSeconds waitForSeconds = new WaitForSeconds(); + waitForSeconds.setTable("test"); + waitForSeconds.setSeconds(5); + waitForSeconds.setCommandExecutionId(commandExecutionId); - try { - TableContextFactory tableContextFactory = getTableContextFactory(); - response = _manager.execute(tableContextFactory, "wait", argumentOverlay); - } catch (TimeoutException te) { - _manager.cancel(te.getExecutionId()); - // some how validate the threads have cancelled. - } + Arguments arguments = CommandUtil.toArguments(waitForSeconds, serDe); + BlurObject args = CommandUtil.toBlurObject(arguments); + System.out.println(args.toString(1)); + final ArgumentOverlay argumentOverlay = new ArgumentOverlay(args, serDe); + final AtomicBoolean fail = new AtomicBoolean(); + final AtomicBoolean running = new AtomicBoolean(true); + + new Thread(new Runnable() { + @Override + public void run() { + TableContextFactory tableContextFactory = getTableContextFactory(); + Long instanceExecutionId = null; + while (true) { + try { + Response response; + if (instanceExecutionId == null) { + response = _manager.execute(tableContextFactory, "wait", argumentOverlay); + } else { + response = _manager.reconnect(instanceExecutionId); + } + fail.set(true); + System.out.println(response); + return; + } catch (IOException e) { + if (e.getCause() instanceof CancellationException) { + return; + } + e.printStackTrace(); + fail.set(true); + return; + } catch (TimeoutException e) { + instanceExecutionId = e.getInstanceExecutionId(); + } catch (Exception e) { + e.printStackTrace(); + fail.set(true); + return; + } finally { + running.set(false); + } + } + } + }).start(); + Thread.sleep(1000); + _manager.cancelCommand(commandExecutionId); + Thread.sleep(5000); + if (fail.get() || running.get()) { + fail("Fail [" + fail.get() + "] Running [" + running.get() + "]"); + } } private TableContextFactory getTableContextFactory() {
