[infinispan-commits] Infinispan SVN: r292 - in trunk/core/src/main: java/org/infinispan/config/parsing and 6 other directories.
infinispan-commits at lists.jboss.org
infinispan-commits at lists.jboss.org
Thu May 14 06:55:29 EDT 2009
Author: manik.surtani at jboss.com
Date: 2009-05-14 06:55:29 -0400 (Thu, 14 May 2009)
New Revision: 292
Modified:
trunk/core/src/main/java/org/infinispan/config/Configuration.java
trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java
trunk/core/src/main/java/org/infinispan/factories/KnownComponentNames.java
trunk/core/src/main/java/org/infinispan/factories/NamedExecutorsFactory.java
trunk/core/src/main/java/org/infinispan/interceptors/base/BaseRpcInterceptor.java
trunk/core/src/main/java/org/infinispan/remoting/rpc/CacheRpcManager.java
trunk/core/src/main/java/org/infinispan/remoting/rpc/ResponseMode.java
trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java
trunk/core/src/main/java/org/infinispan/remoting/transport/Transport.java
trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/CommandAwareRpcDispatcher.java
trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java
trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd
Log:
Added separate configuration for asynchronous marshalling
Modified: trunk/core/src/main/java/org/infinispan/config/Configuration.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/config/Configuration.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/config/Configuration.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -213,6 +213,7 @@
private long rehashWaitTime = 60000;
private boolean useLockStriping = true;
private boolean unsafeUnreliableReturnValues = false;
+ private boolean asyncMarshalling = true;
@Start(priority = 1)
private void correctIsolationLevels() {
@@ -463,10 +464,18 @@
this.rehashWaitTime = rehashWaitTime;
}
+ public void setAsyncMarshalling(boolean asyncMarshalling) {
+ testImmutability("asyncMarshalling");
+ this.asyncMarshalling = asyncMarshalling;
+ }
+
// ------------------------------------------------------------------------------------------------------------
// GETTERS
// ------------------------------------------------------------------------------------------------------------
+ public boolean isAsyncMarshalling() {
+ return asyncMarshalling;
+ }
public boolean isUseReplQueue() {
return useReplQueue;
@@ -613,6 +622,7 @@
if (isolationLevel != that.isolationLevel) return false;
if (transactionManagerLookupClass != null ? !transactionManagerLookupClass.equals(that.transactionManagerLookupClass) : that.transactionManagerLookupClass != null)
return false;
+ if (asyncMarshalling != that.asyncMarshalling) return false;
return true;
}
@@ -652,6 +662,7 @@
result = 31 * result + (l1OnRehash ? 1 : 0);
result = 31 * result + (consistentHashClass != null ? consistentHashClass.hashCode() : 0);
result = 31 * result + numOwners;
+ result = 31 * result + (asyncMarshalling ? 1 : 0);
return result;
}
Modified: trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/config/parsing/XmlConfigurationParserImpl.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -347,6 +347,8 @@
if (existsAttribute(tmp)) config.setReplQueueMaxElements(getInt(tmp));
tmp = getAttributeValue(element, "useAsyncSerialization");
if (existsAttribute(tmp)) config.setUseAsyncSerialization(getBoolean(tmp));
+ tmp = getAttributeValue(element, "asyncMarshalling");
+ if (existsAttribute(tmp)) config.setAsyncMarshalling(getBoolean(tmp));
}
void configureLocking(Element element, Configuration config) {
Modified: trunk/core/src/main/java/org/infinispan/factories/KnownComponentNames.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/KnownComponentNames.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/factories/KnownComponentNames.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -8,7 +8,7 @@
* @since 4.0
*/
public class KnownComponentNames {
- public static final String ASYNC_SERIALIZATION_EXECUTOR = "org.infinispan.executors.serialization";
+ public static final String ASYNC_TRANSPORT_EXECUTOR = "org.infinispan.executors.transport";
public static final String ASYNC_NOTIFICATION_EXECUTOR = "org.infinispan.executors.notification";
public static final String EVICTION_SCHEDULED_EXECUTOR = "org.infinispan.executors.eviction";
public static final String ASYNC_REPLICATION_QUEUE_EXECUTOR = "org.infinispan.executors.replicationQueue";
Modified: trunk/core/src/main/java/org/infinispan/factories/NamedExecutorsFactory.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/NamedExecutorsFactory.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/factories/NamedExecutorsFactory.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -26,7 +26,7 @@
if (componentName.equals(KnownComponentNames.ASYNC_NOTIFICATION_EXECUTOR)) {
return (T) buildAndConfigureExecutorService(globalConfiguration.getAsyncListenerExecutorFactoryClass(),
globalConfiguration.getAsyncListenerExecutorProperties());
- } else if (componentName.equals(KnownComponentNames.ASYNC_SERIALIZATION_EXECUTOR)) {
+ } else if (componentName.equals(KnownComponentNames.ASYNC_TRANSPORT_EXECUTOR)) {
return (T) buildAndConfigureExecutorService(globalConfiguration.getAsyncSerializationExecutorFactoryClass(),
globalConfiguration.getAsyncSerializationExecutorProperties());
} else if (componentName.equals(KnownComponentNames.EVICTION_SCHEDULED_EXECUTOR)) {
Modified: trunk/core/src/main/java/org/infinispan/interceptors/base/BaseRpcInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/base/BaseRpcInterceptor.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/interceptors/base/BaseRpcInterceptor.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -50,7 +50,7 @@
@Inject
public void init(CacheRpcManager rpcManager,
- @ComponentName(KnownComponentNames.ASYNC_SERIALIZATION_EXECUTOR) ExecutorService e) {
+ @ComponentName(KnownComponentNames.ASYNC_TRANSPORT_EXECUTOR) ExecutorService e) {
this.rpcManager = rpcManager;
this.asyncExecutorService = e;
}
Modified: trunk/core/src/main/java/org/infinispan/remoting/rpc/CacheRpcManager.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/rpc/CacheRpcManager.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/rpc/CacheRpcManager.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -69,6 +69,10 @@
}
}
+ private ResponseMode getResponseMode(boolean sync) {
+ return sync ? ResponseMode.SYNCHRONOUS : configuration.isAsyncMarshalling() ? ResponseMode.ASYNCHRONOUS : ResponseMode.ASYNCHRONOUS_WITH_SYNC_MARSHALLING;
+ }
+
public void multicastRpcCommand(List<Address> recipients, CacheRpcCommand command, boolean sync, boolean useOutOfBandMessage) throws ReplicationException {
if (trace) {
log.trace("invoking method " + command.getClass().getSimpleName() + ", members=" + rpcManager.getTransport().getMembers() + ", mode=" +
@@ -81,7 +85,7 @@
try {
rsps = rpcManager.invokeRemotely(recipients,
command,
- sync ? ResponseMode.SYNCHRONOUS : ResponseMode.ASYNCHRONOUS, // is synchronised?
+ getResponseMode(sync),
configuration.getSyncReplTimeout(), useOutOfBandMessage, stateTransferEnabled
);
if (trace) log.trace("responses=" + rsps);
Modified: trunk/core/src/main/java/org/infinispan/remoting/rpc/ResponseMode.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/rpc/ResponseMode.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/rpc/ResponseMode.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -7,5 +7,12 @@
* @since 4.0
*/
public enum ResponseMode {
- SYNCHRONOUS, ASYNCHRONOUS, WAIT_FOR_VALID_RESPONSE
+ SYNCHRONOUS,
+ ASYNCHRONOUS,
+ ASYNCHRONOUS_WITH_SYNC_MARSHALLING,
+ WAIT_FOR_VALID_RESPONSE;
+
+ public boolean isAsynchronous() {
+ return this == ASYNCHRONOUS || this == ASYNCHRONOUS_WITH_SYNC_MARSHALLING;
+ }
}
Modified: trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/rpc/RpcManagerImpl.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -49,7 +49,7 @@
@Inject
public void injectDependencies(GlobalConfiguration globalConfiguration, Transport t, InboundInvocationHandler handler,
Marshaller marshaller,
- @ComponentName(KnownComponentNames.ASYNC_SERIALIZATION_EXECUTOR) ExecutorService e,
+ @ComponentName(KnownComponentNames.ASYNC_TRANSPORT_EXECUTOR) ExecutorService e,
CacheManagerNotifier notifier) {
this.t = t;
this.t.initialize(globalConfiguration, globalConfiguration.getTransportProperties(), marshaller, e, handler,
@@ -71,7 +71,7 @@
List<Response> result = t.invokeRemotely(recipients, rpcCommand, mode, timeout, usePriorityQueue, responseFilter, stateTransferEnabled);
if (isStatisticsEnabled()) replicationCount.incrementAndGet();
return result;
- } catch (CacheException e) {
+ } catch (CacheException e) {
if (log.isTraceEnabled()) log.trace("replicaiton exception: ", e);
if (isStatisticsEnabled()) replicationFailures.incrementAndGet();
throw e;
Modified: trunk/core/src/main/java/org/infinispan/remoting/transport/Transport.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/transport/Transport.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/transport/Transport.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -9,9 +9,9 @@
import org.infinispan.marshall.Marshaller;
import org.infinispan.notifications.cachemanagerlistener.CacheManagerNotifier;
import org.infinispan.remoting.InboundInvocationHandler;
+import org.infinispan.remoting.responses.Response;
import org.infinispan.remoting.rpc.ResponseFilter;
import org.infinispan.remoting.rpc.ResponseMode;
-import org.infinispan.remoting.responses.Response;
import org.infinispan.statetransfer.StateTransferException;
import java.util.List;
@@ -59,7 +59,8 @@
* @return a list of responses from each member contacted.
* @throws Exception in the event of problems.
*/
- List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout, boolean usePriorityQueue, ResponseFilter responseFilter, boolean supportReplay) throws Exception;
+ List<Response> invokeRemotely(List<Address> recipients, ReplicableCommand rpcCommand, ResponseMode mode, long timeout,
+ boolean usePriorityQueue, ResponseFilter responseFilter, boolean supportReplay) throws Exception;
/**
* @return true if the current Channel is the coordinator of the cluster.
Modified: trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/CommandAwareRpcDispatcher.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/CommandAwareRpcDispatcher.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/CommandAwareRpcDispatcher.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -96,15 +96,16 @@
* org.jgroups.blocks.RspFilter)} except that this version is aware of {@link ReplicableCommand} objects.
*/
public RspList invokeRemoteCommands(Vector<Address> dests, ReplicableCommand command, int mode, long timeout,
- boolean anycasting, boolean oob, RspFilter filter, boolean supportReplay)
+ boolean anycasting, boolean oob, RspFilter filter, boolean supportReplay, boolean asyncMarshalling)
throws NotSerializableException, ExecutionException, InterruptedException {
+
ReplicationTask task = new ReplicationTask(command, oob, dests, mode, timeout, anycasting, filter, supportReplay);
- if (mode == GroupRequest.GET_NONE) {
+ if (asyncMarshalling) {
asyncExecutor.submit(task);
return null; // don't wait for a response!
} else {
- RspList response = null;
+ RspList response;
try {
response = task.call();
} catch (Exception e) {
Modified: trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/java/org/infinispan/remoting/transport/jgroups/JGroupsTransport.java 2009-05-14 10:55:29 UTC (rev 292)
@@ -8,10 +8,10 @@
import org.infinispan.notifications.cachemanagerlistener.CacheManagerNotifier;
import org.infinispan.remoting.InboundInvocationHandler;
import org.infinispan.remoting.ReplicationException;
+import org.infinispan.remoting.responses.ExceptionResponse;
+import org.infinispan.remoting.responses.Response;
import org.infinispan.remoting.rpc.ResponseFilter;
import org.infinispan.remoting.rpc.ResponseMode;
-import org.infinispan.remoting.responses.ExceptionResponse;
-import org.infinispan.remoting.responses.Response;
import org.infinispan.remoting.transport.Address;
import org.infinispan.remoting.transport.DistributedSync;
import org.infinispan.remoting.transport.Transport;
@@ -289,13 +289,13 @@
boolean unlock = true;
// if there is a FLUSH in progress, block till it completes
flushTracker.blockUntilReleased(distributedSyncTimeout, MILLISECONDS);
-
+ boolean asyncMarshalling = mode == ResponseMode.ASYNCHRONOUS;
try {
RspList rsps = dispatcher.invokeRemoteCommands(toJGroupsAddressVector(recipients), rpcCommand, toJGroupsMode(mode),
timeout, false, usePriorityQueue,
- toJGroupsFilter(responseFilter), supportReplay);
+ toJGroupsFilter(responseFilter), supportReplay, asyncMarshalling);
- if (mode == ResponseMode.ASYNCHRONOUS) return Collections.emptyList();// async case
+ if (mode.isAsynchronous()) return Collections.emptyList();// async case
if (trace)
log.trace("Cache [{0}]: responses for command {1}:\n{2}", getAddress(), rpcCommand.getClass().getSimpleName(), rsps);
@@ -341,6 +341,7 @@
private int toJGroupsMode(ResponseMode mode) {
switch (mode) {
case ASYNCHRONOUS:
+ case ASYNCHRONOUS_WITH_SYNC_MARSHALLING:
return GroupRequest.GET_NONE;
case SYNCHRONOUS:
return GroupRequest.GET_ALL;
Modified: trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd
===================================================================
--- trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd 2009-05-14 10:49:52 UTC (rev 291)
+++ trunk/core/src/main/resources/schema/infinispan-config-4.0.xsd 2009-05-14 10:55:29 UTC (rev 292)
@@ -19,7 +19,7 @@
<xs:element name="globalJmxStatistics" type="tns:globalJmxStatisticsType" minOccurs="0" maxOccurs="1"/>
<xs:element name="transport" type="tns:transportConfigType" minOccurs="0" maxOccurs="1"/>
<xs:element name="asyncListenerExecutor" type="tns:executorConfigurationType" minOccurs="0" maxOccurs="1"/>
- <xs:element name="asyncSerializationExecutor" type="tns:executorConfigurationType" minOccurs="0"
+ <xs:element name="asyncTransportExecutor" type="tns:executorConfigurationType" minOccurs="0"
maxOccurs="1"/>
<xs:element name="evictionScheduledExecutor" type="tns:executorConfigurationType" minOccurs="0"
maxOccurs="1"/>
@@ -176,6 +176,7 @@
</xs:complexType>
<xs:complexType name="asyncType">
+ <xs:attribute name="asyncMarshalling" type="tns:booleanType"/>
<xs:attribute name="useReplQueue" type="tns:booleanType"/>
<xs:attribute name="replQueueInterval" type="tns:positiveNumber"/>
<xs:attribute name="replQueueMaxElements" type="tns:positiveNumber"/>
More information about the infinispan-commits
mailing list