[infinispan-commits] Infinispan SVN: r1696 - trunk/client/hotrod-client/src/main/java/org/infinispan/client/hotrod/impl/transport/tcp.
infinispan-commits at lists.jboss.org
infinispan-commits at lists.jboss.org
Thu Apr 15 19:46:26 EDT 2010
Author: mircea.markus
Date: 2010-04-15 19:46:26 -0400 (Thu, 15 Apr 2010)
New Revision: 1696
Modified:
trunk/client/hotrod-client/src/main/java/org/infinispan/client/hotrod/impl/transport/tcp/RoundRobinBalancingStrategy.java
trunk/client/hotrod-client/src/main/java/org/infinispan/client/hotrod/impl/transport/tcp/TcpTransportFactory.java
Log:
better sync
Modified: trunk/client/hotrod-client/src/main/java/org/infinispan/client/hotrod/impl/transport/tcp/RoundRobinBalancingStrategy.java
===================================================================
--- trunk/client/hotrod-client/src/main/java/org/infinispan/client/hotrod/impl/transport/tcp/RoundRobinBalancingStrategy.java 2010-04-15 23:22:48 UTC (rev 1695)
+++ trunk/client/hotrod-client/src/main/java/org/infinispan/client/hotrod/impl/transport/tcp/RoundRobinBalancingStrategy.java 2010-04-15 23:46:26 UTC (rev 1696)
@@ -1,8 +1,15 @@
package org.infinispan.client.hotrod.impl.transport.tcp;
+import net.jcip.annotations.ThreadSafe;
+import org.infinispan.util.logging.Log;
+import org.infinispan.util.logging.LogFactory;
+
import java.net.InetSocketAddress;
import java.util.Collection;
import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReadWriteLock;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
/**
* // TODO: Document this
@@ -12,19 +19,43 @@
* @author Mircea.Markus at jboss.com
* @since 4.1
*/
+ at ThreadSafe
public class RoundRobinBalancingStrategy implements RequestBalancingStrategy {
+ private static Log log = LogFactory.getLog(RoundRobinBalancingStrategy.class);
+
+ private final ReadWriteLock readWriteLock = new ReentrantReadWriteLock();
+ private final Lock readLock = readWriteLock.readLock();
+ private final Lock writeLock = readWriteLock.writeLock();
+ private final AtomicInteger index = new AtomicInteger();
+
private InetSocketAddress[] servers;
- private AtomicInteger index = new AtomicInteger();
@Override
public void setServers(Collection<InetSocketAddress> servers) {
- this.servers = servers.toArray(new InetSocketAddress[servers.size()]);
+ writeLock.lock();
+ try {
+ this.servers = servers.toArray(new InetSocketAddress[servers.size()]);
+ } finally {
+ writeLock.unlock();
+ }
}
+ /**
+ * Multiple threads might call this method at the same time.
+ */
@Override
public InetSocketAddress nextServer() {
- int pos = index.incrementAndGet() % servers.length;
- return servers[pos];
+ readLock.lock();
+ try {
+ int pos = index.incrementAndGet() % servers.length;
+ InetSocketAddress server = servers[pos];
+ if (log.isTraceEnabled()) {
+ log.trace("Retuning server: " + server);
+ }
+ return server;
+ } finally {
+ readLock.unlock();
+ }
}
}
Modified: trunk/client/hotrod-client/src/main/java/org/infinispan/client/hotrod/impl/transport/tcp/TcpTransportFactory.java
===================================================================
--- trunk/client/hotrod-client/src/main/java/org/infinispan/client/hotrod/impl/transport/tcp/TcpTransportFactory.java 2010-04-15 23:22:48 UTC (rev 1695)
+++ trunk/client/hotrod-client/src/main/java/org/infinispan/client/hotrod/impl/transport/tcp/TcpTransportFactory.java 2010-04-15 23:46:26 UTC (rev 1696)
@@ -25,7 +25,7 @@
private static Log log = LogFactory.getLog(TcpTransportFactory.class);
- private GenericKeyedObjectPool connectionPool;
+ private volatile GenericKeyedObjectPool connectionPool;
private PropsKeyedObjectPoolFactory poolFactory;
private RequestBalancingStrategy balancer;
private Collection<InetSocketAddress> servers;
More information about the infinispan-commits
mailing list