CVS update: JGroups/src/org/jgroups/protocols UFC.java MFC.java FlowControl.java
"Bela Ban" <[email protected]>
| Newsgroups | gmane.comp.java.javagroups.cvs |
|---|---|
| Message-ID | <[email protected]> |
User: belaban
Date: 10/09/09 11:34:47
Modified: src/org/jgroups/protocols UFC.java MFC.java FlowControl.java
Log:
first draft of separate impls for unicast (UFC) and multicast (MFC) flow control (https://jira.jboss.org/browse/JGRP-1154)
Revision Changes Path
1.3 +64 -99 JGroups/src/org/jgroups/protocols/UFC.java
Index: UFC.java
===================================================================
RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/UFC.java,v
retrieving revision 1.2
retrieving revision 1.3
diff -u -r1.2 -r1.3
--- UFC.java 7 Sep 2010 10:38:21 -0000 1.2
+++ UFC.java 9 Sep 2010 11:34:47 -0000 1.3
@@ -4,6 +4,13 @@
import org.jgroups.Event;
import org.jgroups.Message;
import org.jgroups.annotations.MBean;
+import org.jgroups.annotations.ManagedAttribute;
+import org.jgroups.annotations.ManagedOperation;
+
+import java.util.Iterator;
+import java.util.Map;
+import java.util.Vector;
+import java.util.concurrent.ConcurrentHashMap;
/**
@@ -24,50 +31,62 @@
* <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
* <ol/>
* @author Bela Ban
- * @version $Id: UFC.java,v 1.2 2010/09/07 10:38:21 belaban Exp $
+ * @version $Id: UFC.java,v 1.3 2010/09/09 11:34:47 belaban Exp $
*/
@MBean(description="Simple flow control protocol based on a credit system")
public class UFC extends FlowControl {
-
+ /**
+ * Map<Address,Long>: keys are members, values are credits left. For each send,
+ * the number of credits is decremented by the message size
+ */
+ protected final Map<Address,Credit> sent=new ConcurrentHashMap<Address,Credit>(11);
- protected boolean handleMulticastMessage() {
- return false;
- }
- protected Credit createCredit(long credits) {
- return new UfcCredit(credits);
+ @ManagedOperation(description="Print sender credits")
+ public String printSenderCredits() {
+ return printMap(sent);
}
- public void unblock() {
- super.unblock();
+
+ @ManagedOperation(description="Print credits")
+ public String printCredits() {
+ StringBuilder sb=new StringBuilder(super.printCredits());
+ sb.append("\nsenders:\n").append(printMap(sent));
+ return sb.toString();
}
- public double getAverageTimeBlocked() {
- int blockings=0;
- long total_time_blocked=0;
+ public Map<String, Object> dumpStats() {
+ Map<String, Object> retval=super.dumpStats();
+ retval.put("senders", printMap(sent));
+ return retval;
+ }
- for(Credit cred: sent.values()) {
- blockings+=((UfcCredit)cred).getNumBlockings();
- total_time_blocked=((UfcCredit)cred).getTotalBlockingTime();
- }
- return blockings > 0? total_time_blocked / (double)blockings : 0.0; // prevent div-by-zero
+ protected boolean handleMulticastMessage() {
+ return false;
}
+
+ public void unblock() {
+ super.unblock();
+ }
+
+ @ManagedAttribute(description="Number of times flow control blocks sender")
public int getNumberOfBlockings() {
int retval=0;
for(Credit cred: sent.values())
- retval+=((UfcCredit)cred).getNumBlockings();
+ retval+=cred.getNumBlockings();
return retval;
}
+ @ManagedAttribute(description="Total time (ms) spent in flow control block")
public long getTotalTimeBlocked() {
long retval=0;
for(Credit cred: sent.values())
- retval+=((UfcCredit)cred).getTotalBlockingTime();
+ retval+=cred.getTotalBlockingTime();
return retval;
}
@@ -91,13 +110,7 @@
return down_prot.down(evt);
}
- if(ignore_synchronous_response && ignore_thread.get()) { // JGRP-465
- if(log.isTraceEnabled())
- log.trace("bypassing blocking to avoid deadlocking " + Thread.currentThread());
- return down_prot.down(evt);
- }
-
- UfcCredit cred=(UfcCredit)sent.get(dest);
+ Credit cred=sent.get(dest);
if(cred == null) {
log.error("destination " + dest + " not found; passing message down");
return down_prot.down(evt);
@@ -106,36 +119,47 @@
long block_time=max_block_times != null? getMaxBlockTime(length) : max_block_time;
while(running) {
- if(cred.decrementIfEnoughCredits(length, 0)) // timeout == 0: don't block
- break;
-
- if(log.isTraceEnabled())
- log.trace("blocking for credits (for " + block_time + " ms)");
boolean rc=cred.decrementIfEnoughCredits(length, block_time);
- if(rc && log.isTraceEnabled())
- log.trace("unblocking (received credits)");
-
if(rc || !running || max_block_times != null)
break;
- if(cred.sendCreditRequest(System.currentTimeMillis()))
- sendCreditRequest(dest, cred.get());
+ if(cred.needToSendCreditRequest())
+ sendCreditRequest(dest, Math.max(0, max_credits - cred.get()));
}
// send message - either after regular processing, or after blocking (when enough credits available again)
return down_prot.down(evt);
}
-
- protected void handleCredit(Address sender, Number increase) {
+ protected void handleViewChange(Vector<Address> mbrs) {
+ super.handleViewChange(mbrs);
+ if(mbrs == null) return;
+
+ // add members not in membership to received and sent hashmap (with full credits)
+ for(int i=0; i < mbrs.size(); i++) {
+ Address addr=mbrs.elementAt(i);
+ if(!sent.containsKey(addr))
+ sent.put(addr, new Credit(max_credits));
+ }
+
+ // remove members that left
+ for(Iterator<Address> it=sent.keySet().iterator(); it.hasNext();) {
+ Address addr=it.next();
+ if(!mbrs.contains(addr))
+ it.remove(); // modified the underlying map
+ }
+ }
+
+
+ protected void handleCredit(Address sender, long increase) {
if(sender == null) return;
StringBuilder sb=null;
Credit cred=sent.get(sender);
if(cred == null)
return;
- long new_credit=Math.min(max_credits, cred.get() + increase.longValue());
+ long new_credit=Math.min(max_credits, cred.get() + increase);
if(log.isTraceEnabled()) {
sb=new StringBuilder();
@@ -143,67 +167,8 @@
.append(", new credits: ").append(new_credit);
log.trace(sb);
}
-
- cred.increment(increase.longValue());
+ cred.increment(increase);
}
- protected class UfcCredit extends Credit {
- int num_blockings=0;
- long total_blocking_time=0;
- long last_credit_request=0;
-
- protected UfcCredit(long credits) {
- super(credits);
- }
-
- protected synchronized boolean decrementIfEnoughCredits(long credits, long timeout) {
- if(credits <= credits_left) {
- credits_left-=credits;
- return true;
- }
-
- if(timeout <= 0)
- return false;
-
- long start=System.currentTimeMillis();
- try {
- this.wait(timeout);
- }
- catch(InterruptedException e) {
- }
- finally {
- total_blocking_time+=System.currentTimeMillis() - start;
- num_blockings++;
- }
-
- if(credits <= credits_left) {
- credits_left-=credits;
- return true;
- }
- return false;
- }
-
- protected synchronized long increment(long credits) {
- long retval=super.increment(credits);
- notifyAll();
- return retval;
- }
-
- protected synchronized boolean sendCreditRequest(long current_time) {
- if(current_time - last_credit_request >= max_block_time) {
- // we have to set this var now, because we release the lock below (for sending a credit request), so
- // all blocked threads would send a credit request, leading to a credit request storm
- last_credit_request=System.currentTimeMillis();
- return true;
- }
- return false;
- }
-
- protected int getNumBlockings() {return num_blockings;}
-
- protected long getTotalBlockingTime() {return total_blocking_time;}
- }
-
-
}
1.4 +81 -195 JGroups/src/org/jgroups/protocols/MFC.java
Index: MFC.java
===================================================================
RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/MFC.java,v
retrieving revision 1.3
retrieving revision 1.4
diff -u -r1.3 -r1.4
--- MFC.java 7 Sep 2010 10:38:32 -0000 1.3
+++ MFC.java 9 Sep 2010 11:34:47 -0000 1.4
@@ -3,15 +3,16 @@
import org.jgroups.Address;
import org.jgroups.Event;
import org.jgroups.Message;
-import org.jgroups.annotations.GuardedBy;
import org.jgroups.annotations.MBean;
+import org.jgroups.annotations.ManagedAttribute;
import org.jgroups.annotations.ManagedOperation;
+import org.jgroups.util.CreditMap;
-import java.util.*;
-import java.util.concurrent.TimeUnit;
-import java.util.concurrent.locks.Condition;
-import java.util.concurrent.locks.Lock;
-import java.util.concurrent.locks.ReentrantLock;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.Vector;
+import java.util.concurrent.atomic.AtomicLong;
/**
@@ -32,7 +33,7 @@
* <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
* <ol/>
* @author Bela Ban
- * @version $Id: MFC.java,v 1.3 2010/09/07 10:38:32 belaban Exp $
+ * @version $Id: MFC.java,v 1.4 2010/09/09 11:34:47 belaban Exp $
*/
@MBean(description="Simple flow control protocol based on a credit system")
public class MFC extends FlowControl {
@@ -41,83 +42,62 @@
/* --------------------------------------------- Fields ------------------------------------------------------ */
-
-
- /**
- * the lowest credits of any destination (sent_msgs)
- */
- @GuardedBy("lock")
- private long lowest_credit=max_credits;
-
- /** Lock protecting sent credits table and some other vars (creditors for example) */
- private final Lock lock=new ReentrantLock();
+ /** Maintains credits per member */
+ protected CreditMap credits;
- /** List of members from whom we expect credits */
- @GuardedBy("lock")
- protected final Set<Address> creditors=new HashSet<Address>(11);
+ /** Number of credits waiting to be decremented, used to determine how many credits to ask for in credit requests */
+ @ManagedAttribute(description="The total number of bytes accumulated by messages which cannot be sent due " +
+ "to insufficient credits",writable=false)
+ protected final AtomicLong blocked_credits=new AtomicLong(0);
+ /** Last time a credit request was sent. Used to prevent credit request storms */
+ protected long last_credit_request=0;
- /** Mutex to block on down() */
- private final Condition credits_available=lock.newCondition();
- /**
- * Allows to unblock a blocked sender from an external program, e.g. JMX
- */
+ /** Allows to unblock a blocked sender from an external program, e.g. JMX */
@ManagedOperation(description="Unblock a sender")
public void unblock() {
- lock.lock();
- try {
- if(log.isTraceEnabled())
- log.trace("unblocking the sender and replenishing all members");
-
- for(Map.Entry<Address,Credit> entry: sent.entrySet())
- entry.getValue().set(max_credits);
+ if(log.isTraceEnabled())
+ log.trace("unblocking the sender and replenishing all members");
+ credits.replenishAll();
+ }
- lowest_credit=computeLowestCredit(sent);
- creditors.clear();
- credits_available.signalAll();
- }
- finally {
- lock.unlock();
- }
+ @ManagedOperation(description="Print credits")
+ public String printCredits() {
+ return super.printCredits() + "\nsenders min credits: " + credits.getMinCredits();
}
-
- public void init() throws Exception {
- super.init();
- lowest_credit=max_credits;
+ @ManagedOperation(description="Print sender credits")
+ public String printSenderCredits() {
+ return credits.toString();
}
- public void start() throws Exception {
- super.start();
- lowest_credit=max_credits;
+ @ManagedAttribute(description="Number of times flow control blocks sender")
+ public int getNumberOfBlockings() {
+ return credits.getNumBlockings();
}
- public void stop() {
- super.stop();
- lock.lock();
- try {
- running=false;
- ignore_thread.set(false);
- credits_available.signalAll(); // notify all threads waiting on the mutex that we are done
- }
- finally {
- lock.unlock();
- }
+ @ManagedAttribute(description="Total time (ms) spent in flow control block")
+ public long getTotalTimeBlocked() {
+ return credits.getTotalBlockTime();
}
protected boolean handleMulticastMessage() {
return true;
}
- protected Credit createCredit(long credits) {
- return new MfcCredit(credits);
+
+ public void init() throws Exception {
+ super.init();
+ credits=new CreditMap(max_credits);
}
-
-
+ public void stop() {
+ super.stop();
+ credits.clear();
+ }
protected Object handleDownMessage(final Event evt, final Message msg, int length) {
Address dest=msg.getDest();
@@ -126,133 +106,54 @@
return down_prot.down(evt);
}
- lock.lock();
+ long block_time=max_block_times != null? getMaxBlockTime(length) : max_block_time;
try {
- if(length > lowest_credit) { // then block and loop asking for credits until enough credits are available
- if(ignore_synchronous_response && ignore_thread.get()) { // JGRP-465
- if(log.isTraceEnabled())
- log.trace("bypassing blocking to avoid deadlocking " + Thread.currentThread());
- }
- else {
- determineCreditors(length);
- num_blockings++; // we count overall blockings, not blockings for *all* threads
- if(log.isTraceEnabled())
- log.trace("Starting blocking. lowest_credit=" + lowest_credit + "; msg length =" + length);
-
- long block_time=max_block_times != null? getMaxBlockTime(length) : max_block_time;
- while(length > lowest_credit && running) {
- try {
- long start=System.currentTimeMillis();
- boolean rc=credits_available.await(block_time, TimeUnit.MILLISECONDS);
- total_time_blocking+=System.currentTimeMillis() - start;
- if(length <= lowest_credit || rc || !running)
- break;
-
- // if we use max_block_times, then we do *not* send credit requests, even if we run
- // into timeouts: in this case, it is up to the receivers to send new credits
- if(!rc && max_block_times != null)
- break;
-
- long curr_time=System.currentTimeMillis();
- long wait_time=curr_time - last_credit_request;
- if(wait_time >= max_block_time) {
-
- // we have to set this var now, because we release the lock below (for sending a
- // credit request), so all blocked threads would send a credit request, leading to
- // a credit request storm
- last_credit_request=curr_time;
-
- // we need to send the credit requests down *without* holding the lock, otherwise we might
- // run into the deadlock described in http://jira.jboss.com/jira/browse/JGRP-292
- Map<Address,Credit> sent_copy=new HashMap<Address,Credit>(sent);
- sent_copy.keySet().retainAll(creditors);
- lock.unlock();
- try {
- for(Map.Entry<Address,Credit> entry: sent_copy.entrySet())
- sendCreditRequest(entry.getKey(), entry.getValue().get());
- }
- finally {
- lock.lock();
- }
- }
- }
- catch(InterruptedException e) {
- // bela June 15 2007: don't interrupt the thread again, as this will trigger an infinite loop !!
- // (http://jira.jboss.com/jira/browse/JGRP-536)
- // Thread.currentThread().interrupt();
- }
- }
+ if(length > 0 && max_block_times == null)
+ this.blocked_credits.addAndGet(length);
+ while(running) {
+ boolean rc=credits.decrement(length, block_time);
+ if(rc || max_block_times != null || !running)
+ break;
+
+ if(needToSendCreditRequest()) {
+ long credits_blocked=blocked_credits.get();
+ List<Address> targets=credits.getMembersWithInsufficientCredits(credits_blocked);
+ for(Address target: targets)
+ sendCreditRequest(target, credits_blocked);
}
}
-
- long tmp=decrementCredit(sent, length);
- if(tmp != -1)
- lowest_credit=Math.min(tmp, lowest_credit);
}
finally {
- lock.unlock();
+ if(length > 0 && max_block_times == null)
+ this.blocked_credits.getAndAdd(-length);
}
- // send message - either after regular processing, or after blocking (when enough credits available again)
+ // send message - either after regular processing, or after blocking (when enough credits are available again)
return down_prot.down(evt);
}
- /**
- * Adds members which have not enough credits to the creditors list. Called with lock held
- * @param length
- */
- protected void determineCreditors(int length) {
- for(Map.Entry<Address,Credit> entry: sent.entrySet()) {
- if(entry.getValue().get() <= length)
- creditors.add(entry.getKey());
- }
- }
-
-
- /**
- * Decrements credits from a single member, or all members in sent_msgs, depending on whether it is a multicast
- * or unicast message. No need to acquire mutex (must already be held when this method is called)
- * @param dest
- * @param credits
- * @return The lowest number of credits left, or -1 if a unicast member was not found
- */
- protected long decrementCredit(Map<Address,Credit> map, long credits) {
- if(map.isEmpty())
- return -1;
- long lowest=max_credits;
- for(Credit cred: map.values())
- lowest=Math.min(cred.decrement(credits), lowest);
- return lowest;
- }
-
- protected void handleCredit(Address sender, Number increase) {
- if(sender == null) return;
- StringBuilder sb=null;
- lock.lock();
- try {
- Credit cred=sent.get(sender);
- if(cred == null)
- return;
- long new_credit=Math.min(max_credits, cred.get() + increase.longValue());
-
- if(log.isTraceEnabled()) {
- sb=new StringBuilder();
- sb.append("received " + increase + " credits from ").append(sender).append(", old credits: ").append(cred)
- .append(", new credits: ").append(new_credit).append(".\nCreditors before are: ").append(creditors);
- log.trace(sb);
- }
+ protected synchronized boolean needToSendCreditRequest() {
+ long curr_time=System.currentTimeMillis();
+ long wait_time=curr_time - last_credit_request;
+ if(wait_time >= max_block_time) {
+ last_credit_request=curr_time;
+ return true;
+ }
+ return false;
+ }
+
- cred.increment(increase.longValue());
- lowest_credit=computeLowestCredit(sent);
- if(!creditors.isEmpty() && creditors.remove(sender) && creditors.isEmpty())
- credits_available.signalAll();
- }
- finally {
- lock.unlock();
+ protected void handleCredit(Address sender, long increase) {
+ credits.replenish(sender, increase);
+ if(log.isTraceEnabled()) {
+ StringBuilder sb=new StringBuilder();
+ sb.append("received " + increase + " credits from ").append(sender).append(", new credits for " + sender + " : ")
+ .append(credits.get(sender) + ", min_credits=" + credits.getMinCredits());
+ log.trace(sb);
}
}
@@ -260,32 +161,17 @@
protected void handleViewChange(Vector<Address> mbrs) {
super.handleViewChange(mbrs);
- lock.lock();
- try {
- // fixed http://jira.jboss.com/jira/browse/JGRP-754 (CCME)
- for(Iterator<Address> it=creditors.iterator(); it.hasNext();) {
- Address creditor=it.next();
- if(!mbrs.contains(creditor))
- it.remove();
- }
-
- if(log.isTraceEnabled()) log.trace("creditors are " + creditors);
- if(creditors.isEmpty()) {
- lowest_credit=computeLowestCredit(sent);
- credits_available.signalAll();
- }
- }
- finally {
- lock.unlock();
+ Set<Address> keys=new HashSet<Address>(credits.keys());
+ for(Address key: keys) {
+ if(!mbrs.contains(key))
+ credits.remove(key);
}
+
+ for(Address key: mbrs)
+ credits.putIfAbsent(key);
}
- protected class MfcCredit extends Credit {
- protected MfcCredit(long credits) {
- super(credits);
- }
- }
}
1.5 +99 -173 JGroups/src/org/jgroups/protocols/FlowControl.java
Index: FlowControl.java
===================================================================
RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/FlowControl.java,v
retrieving revision 1.4
retrieving revision 1.5
diff -u -r1.4 -r1.5
--- FlowControl.java 7 Sep 2010 10:38:01 -0000 1.4
+++ FlowControl.java 9 Sep 2010 11:34:47 -0000 1.5
@@ -6,7 +6,6 @@
import org.jgroups.View;
import org.jgroups.annotations.*;
import org.jgroups.stack.Protocol;
-import org.jgroups.util.BoundedList;
import org.jgroups.util.Util;
import java.util.*;
@@ -31,7 +30,7 @@
* <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
* <ol/>
* @author Bela Ban
- * @version $Id: FlowControl.java,v 1.4 2010/09/07 10:38:01 belaban Exp $
+ * @version $Id: FlowControl.java,v 1.5 2010/09/09 11:34:47 belaban Exp $
*/
@MBean(description="Simple flow control protocol based on a credit system")
public abstract class FlowControl extends Protocol {
@@ -88,10 +87,9 @@
protected long min_credits=0;
/**
- * Whether an up thread that comes back down should be allowed to
- * bypass blocking if all credits are exhausted. Avoids JGRP-465.
- * Set to false by default in 2.5 because we have OOB messages for credit replenishments - this flag should not be set
- * to true if the concurrent stack is used
+ * Whether an up thread that comes back down should be allowed to bypass blocking if all credits are exhausted.
+ * Avoids JGRP-465. Set to false by default in 2.5 because we have OOB messages for credit replenishments -
+ * this flag should not be set to true if the concurrent stack is used
*/
@Property(description="Does not block a down message if it is a result of handling an up message in the" +
"same thread. Fixes JGRP-928")
@@ -101,27 +99,12 @@
/* --------------------------------------------- JMX ------------------------------------------------------ */
-
-
- protected int num_blockings=0;
- protected int num_credit_requests_received=0, num_credit_requests_sent=0;
- protected int num_credit_responses_sent=0, num_credit_responses_received=0;
- protected long total_time_blocking=0;
+ protected int num_credit_requests_received=0, num_credit_requests_sent=0;
+ protected int num_credit_responses_sent=0, num_credit_responses_received=0;
+
- // protected final BoundedList<Long> last_blockings=new BoundedList<Long>(50);
-
-
-
/* --------------------------------------------- Fields ------------------------------------------------------ */
-
-
- /**
- * Map<Address,Long>: keys are members, values are credits left. For each send, the
- * number of credits is decremented by the message size. A HashMap rather than a ConcurrentHashMap is
- * currently used as there might be null values
- */
- @GuardedBy("lock")
- protected final Map<Address,Credit> sent=new ConcurrentHashMap<Address,Credit>(11);
+
/**
* Keeps track of credits / member at the receiver's side. Keys are members, values are credits left (in bytes).
@@ -132,11 +115,6 @@
protected final Map<Address,Credit> received=new ConcurrentHashMap<Address,Credit>(11);
-
-
- /** Peers who have asked for credit that we didn't have */
- protected final Set<Address> pending_requesters=new HashSet<Address>(11);
-
/**
* Whether FlowControl is still running, this is set to false when the protocol terminates (on stop())
*/
@@ -159,16 +137,10 @@
}
};
- /** Last time a credit request was sent. Used to prevent credit request storms */
- @GuardedBy("lock")
- protected long last_credit_request=0;
public void resetStats() {
super.resetStats();
- num_blockings=0;
num_credit_responses_sent=num_credit_responses_received=num_credit_requests_received=num_credit_requests_sent=0;
- total_time_blocking=0;
- // last_blockings.clear();
}
public long getMaxCredits() {
@@ -195,10 +167,8 @@
this.min_credits=min_credits;
}
- @ManagedAttribute(description="Number of times flow control blocks sender")
- public int getNumberOfBlockings() {
- return num_blockings;
- }
+
+ public abstract int getNumberOfBlockings();
public long getMaxBlockTime() {
return max_block_time;
@@ -259,14 +229,13 @@
return sb.toString();
}
- @ManagedAttribute(description="Total time (ms) spent in flow control block")
- public long getTotalTimeBlocked() {
- return total_time_blocking;
- }
+
+ public abstract long getTotalTimeBlocked();
@ManagedAttribute(description="Average time spent in a flow control block")
public double getAverageTimeBlocked() {
- return num_blockings == 0? 0.0 : total_time_blocking / (double)num_blockings;
+ long number_of_blockings=getNumberOfBlockings();
+ return number_of_blockings == 0? 0.0 : getTotalTimeBlocked() / (double)number_of_blockings;
}
@ManagedAttribute(description="Number of credit requests received")
@@ -288,36 +257,27 @@
public int getNumberOfCreditResponsesSent() {
return num_credit_responses_sent;
}
-
- @ManagedOperation(description="Print sender credits")
- public String printSenderCredits() {
- return printMap(sent);
- }
+
+ public abstract String printSenderCredits();
@ManagedOperation(description="Print receiver credits")
public String printReceiverCredits() {
return printMap(received);
}
- @ManagedOperation(description="Print credits")
+
public String printCredits() {
StringBuilder sb=new StringBuilder();
- sb.append("senders:\n").append(printMap(sent)).append("\n\nreceivers:\n").append(printMap(received));
+ sb.append("receivers:\n").append(printMap(received));
return sb.toString();
}
public Map<String, Object> dumpStats() {
Map<String, Object> retval=super.dumpStats();
- retval.put("senders", printMap(sent));
retval.put("receivers", printMap(received));
return retval;
}
- //@ManagedOperation(description="Print last blocking times")
- //public String showLastBlockingTimes() {
- // return last_blockings.toString();
- //}
-
protected long getMaxBlockTime(long length) {
if(max_block_times == null)
@@ -339,9 +299,7 @@
*/
protected abstract boolean handleMulticastMessage();
- protected abstract Credit createCredit(long credits);
-
- protected abstract void handleCredit(Address sender, Number increase);
+ protected abstract void handleCredit(Address sender, long increase);
/**
@@ -394,6 +352,13 @@
int length=msg.getLength();
if(length == 0)
break;
+
+ if(ignore_synchronous_response && ignore_thread.get()) { // JGRP-465
+ if(log.isTraceEnabled())
+ log.trace("bypassing flow control because of synchronous response " + Thread.currentThread());
+ break;
+ }
+
return handleDownMessage(evt, msg, length);
case Event.CONFIG:
handleConfigEvent((Map<String,Object>)evt.getArg());
@@ -430,14 +395,14 @@
switch(hdr.type) {
case FcHeader.REPLENISH:
num_credit_responses_received++;
- handleCredit(msg.getSrc(), (Number)msg.getObject());
+ handleCredit(msg.getSrc(), (Long)msg.getObject());
break;
case FcHeader.CREDIT_REQUEST:
num_credit_requests_received++;
Address sender=msg.getSrc();
- Long sent_credits=(Long)msg.getObject();
- if(sent_credits != null)
- handleCreditRequest(received, sender, sent_credits.longValue());
+ Long requested_credits=(Long)msg.getObject();
+ if(requested_credits != null)
+ handleCreditRequest(received, sender, requested_credits.longValue());
break;
default:
log.error("header type " + hdr.type + " not known");
@@ -459,10 +424,8 @@
finally {
if(ignore_synchronous_response)
ignore_thread.set(false); // need to revert because the thread is placed back into the pool
- if(new_credits > 0) {
- if(log.isTraceEnabled()) log.trace("sending " + new_credits + " credits to " + sender);
+ if(new_credits > 0)
sendCredit(sender, new_credits);
- }
}
case Event.VIEW_CHANGE:
@@ -498,39 +461,6 @@
-
-
-
- /**
- * Decrements credits from a single member, or all members in sent_msgs, depending on whether it is a multicast
- * or unicast message. No need to acquire mutex (must already be held when this method is called)
- * @param dest
- * @param credits
- * @return The lowest number of credits left, or -1 if a unicast member was not found
- */
- protected long decrementCredit(Map<Address,Credit> map, Address dest, long credits) {
- if(dest == null || dest.isMulticastAddress()) {
- if(map.isEmpty())
- return -1;
- long lowest=max_credits;
- for(Credit cred: map.values())
- lowest=Math.min(cred.decrement(credits), lowest);
- return lowest;
- }
- else {
- Credit cred=map.get(dest);
- if(cred != null)
- return cred.decrement(credits);
- }
- return -1;
- }
-
-
-
-
-
-
-
/**
* Check whether sender has enough credits left. If not, send it some more
* @param map The hashmap to use
@@ -548,7 +478,7 @@
return 0;
if(log.isTraceEnabled())
- log.trace("sender " + sender + " minus " + length + " credits, " + (cred.get() - length) + " remaining");
+ log.trace(sender + " used " + length + " credits, " + (cred.get() - length) + " remaining");
return cred.decrementAndGet(length);
}
@@ -556,65 +486,27 @@
/**
* @param map The map to modify
* @param sender The sender who requests credits
- * @param left_credits Number of bytes that the sender has left to send messages to us
+ * @param requested_credits Number of bytes that the sender has left to send messages to us
*/
- protected void handleCreditRequest(Map<Address,Credit> map, Address sender, long left_credits) {
- if(sender == null) return;
- long credit_response=0;
- Credit cred=map.get(sender);
-
- long old_credit=cred != null? cred.get() : 0;
- if(old_credit > 0)
- credit_response=Math.min(max_credits, max_credits - old_credit);
+ protected void handleCreditRequest(Map<Address,Credit> map, Address sender, long requested_credits) {
+ Credit cred;
+ if(sender == null || (cred=map.get(sender)) == null)
+ return;
+ long credit_response=Math.min(max_credits, Math.min(requested_credits, max_credits - cred.get()));
if(credit_response > 0) {
if(log.isTraceEnabled())
log.trace("received credit request from " + sender + ": sending " + credit_response + " credits");
- if(cred != null)
- cred.set(max_credits);
- else
- map.put(sender, createCredit(max_credits));
- pending_requesters.remove(sender);
- }
- else {
- if(pending_requesters.contains(sender)) {
- // a sender might have negative credits, e.g. -20000. If we subtracted -20000 from max_credits,
- // we'd end up with max_credits + 20000, and send too many credits back. So if the sender's
- // credits is negative, we simply send max_credits back
- long credits_left=Math.max(0, left_credits);
- credit_response=max_credits - credits_left;
- // credit_response = max_credits;
- if(cred != null)
- cred.set(max_credits);
- else
- map.put(sender, createCredit(max_credits));
- pending_requesters.remove(sender);
- if(log.isWarnEnabled())
- log.warn("Received two credit requests from " + sender +
- " without any intervening messages; sending " + credit_response + " credits");
- }
- else {
- pending_requesters.add(sender);
- if(log.isTraceEnabled())
- log.trace("received credit request from " + sender + " but have no credits available");
- }
- }
-
-
- if(credit_response > 0)
+ cred.set(max_credits);
sendCredit(sender, credit_response);
+ }
}
- protected void sendCredit(Address dest, long credit) {
+ protected void sendCredit(Address dest, long credits) {
if(log.isTraceEnabled())
- log.trace("replenishing " + dest + " with " + credit + " credits");
- Number number;
- if(credit < Integer.MAX_VALUE)
- number=(int)credit;
- else
- number=credit;
- Message msg=new Message(dest, null, number);
+ if(log.isTraceEnabled()) log.trace("sending " + credits + " credits to " + dest);
+ Message msg=new Message(dest, null, new Long(credits));
msg.setFlag(Message.OOB);
msg.putHeader(this.id, REPLENISH_HDR);
down_prot.down(new Event(Event.MSG, msg));
@@ -625,12 +517,12 @@
* We cannot send this request as OOB messages, as the credit request needs to queue up behind the regular messages;
* if a receiver cannot process the regular messages, that is a sign that the sender should be throttled !
* @param dest The member to which we send the credit request
- * @param credits_left The number of bytes (of credits) left for dest
+ * @param credits_needed The number of bytes (of credits) left for dest
*/
- protected void sendCreditRequest(final Address dest, Long credits_left) {
+ protected void sendCreditRequest(final Address dest, Long credits_needed) {
if(log.isTraceEnabled())
log.trace("sending credit request to " + dest);
- Message msg=new Message(dest, null, credits_left);
+ Message msg=new Message(dest, null, credits_needed);
msg.putHeader(this.id, CREDIT_REQUEST_HDR);
down_prot.down(new Event(Event.MSG, msg));
num_credit_requests_sent++;
@@ -647,9 +539,7 @@
for(int i=0; i < mbrs.size(); i++) {
addr=mbrs.elementAt(i);
if(!received.containsKey(addr))
- received.put(addr, createCredit(max_credits));
- if(!sent.containsKey(addr))
- sent.put(addr, createCredit(max_credits));
+ received.put(addr, new Credit(max_credits));
}
// remove members that left
for(Iterator<Address> it=received.keySet().iterator(); it.hasNext();) {
@@ -657,17 +547,9 @@
if(!mbrs.contains(addr))
it.remove();
}
-
- // remove members that left
- for(Iterator<Address> it=sent.keySet().iterator(); it.hasNext();) {
- addr=it.next();
- if(!mbrs.contains(addr))
- it.remove(); // modified the underlying map
- }
-
-
}
+
protected static long computeLowestCredit(Map<Address,Credit> m) {
Collection<Credit> credits=m.values();
return Collections.min(credits).get();
@@ -683,13 +565,46 @@
- protected abstract class Credit implements Comparable {
+ protected class Credit implements Comparable {
protected long credits_left;
+ protected int num_blockings=0;
+ protected long total_blocking_time=0;
+ protected long last_credit_request=0;
+
protected Credit(long credits) {
this.credits_left=credits;
}
+
+ protected synchronized boolean decrementIfEnoughCredits(long credits, long timeout) {
+ if(credits <= credits_left) {
+ credits_left-=credits;
+ return true;
+ }
+
+ if(timeout <= 0)
+ return false;
+
+ long start=System.currentTimeMillis();
+ try {
+ this.wait(timeout);
+ }
+ catch(InterruptedException e) {
+ }
+ finally {
+ total_blocking_time+=System.currentTimeMillis() - start;
+ num_blockings++;
+ }
+
+ if(credits <= credits_left) {
+ credits_left-=credits;
+ return true;
+ }
+ return false;
+ }
+
+
protected synchronized long decrementAndGet(long credits) {
credits_left=Math.max(0, credits_left - credits);
long credit_response=max_credits - credits_left;
@@ -700,19 +615,30 @@
return 0;
}
- protected synchronized long decrement(long credits) {
- return credits_left=Math.max(0, credits_left - credits);
+
+ protected synchronized long increment(long credits) {
+ long retval=credits_left=Math.min(max_credits, credits_left + credits);
+ notifyAll();
+ return retval;
}
-
+ protected synchronized boolean needToSendCreditRequest() {
+ long current_time=System.currentTimeMillis();
+ if(current_time - last_credit_request >= max_block_time) {
+ last_credit_request=current_time;
+ return true;
+ }
+ return false;
+ }
+
+ protected int getNumBlockings() {return num_blockings;}
+
+ protected long getTotalBlockingTime() {return total_blocking_time;}
protected synchronized long get() {return credits_left;}
protected synchronized void set(long new_credits) {credits_left=Math.min(max_credits, new_credits);}
- protected synchronized long increment(long credits) {
- return credits_left=Math.min(max_credits, credits_left + credits);
- }
public String toString() {
return String.valueOf(credits_left);
------------------------------------------------------------------------------
This SF.net Dev2Dev email is sponsored by:
Show off your parallel programming skills.
Enter the Intel(R) Threading Challenge 2010.
http://p.sf.net/sfu/intel-thread-sfd