CVS update: JGroups/src/org/jgroups/protocols DAISYCHAIN.java

"Bela Ban" <[email protected]> Tue, 14 Sep 2010 11:55:36 +0000
Newsgroups gmane.comp.java.javagroups.cvs
Message-ID <[email protected]>
  User: belaban 
  Date: 10/09/14 11:55:36

  Modified:    src/org/jgroups/protocols DAISYCHAIN.java
  Log:
  ns
  
  Revision  Changes    Path
  1.12      +63 -19    JGroups/src/org/jgroups/protocols/DAISYCHAIN.java
  
  Index: DAISYCHAIN.java
  ===================================================================
  RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/DAISYCHAIN.java,v
  retrieving revision 1.11
  retrieving revision 1.12
  diff -u -r1.11 -r1.12
  --- DAISYCHAIN.java	13 Sep 2010 12:02:47 -0000	1.11
  +++ DAISYCHAIN.java	14 Sep 2010 11:55:36 -0000	1.12
  @@ -3,11 +3,13 @@
   import org.jgroups.*;
   import org.jgroups.annotations.*;
   import org.jgroups.stack.Protocol;
  +import org.jgroups.util.ConcurrentLinkedBlockingQueue;
   import org.jgroups.util.Util;
   
   import java.io.DataInputStream;
   import java.io.DataOutputStream;
   import java.io.IOException;
  +import java.util.concurrent.BlockingQueue;
   import java.util.concurrent.Executor;
   
   /**
  @@ -18,7 +20,7 @@
    * send another message. This leads to much better throughput, see the ref in the JIRA.<p/> 
    * JIRA: https://jira.jboss.org/browse/JGRP-1021
    * @author Bela Ban
  - * @version $Id: DAISYCHAIN.java,v 1.11 2010/09/13 12:02:47 belaban Exp $
  + * @version $Id: DAISYCHAIN.java,v 1.12 2010/09/14 11:55:36 belaban Exp $
    */
   @Experimental @Unsupported
   @MBean(description="Protocol just above the transport which disseminates multicasts via daisy chaining")
  @@ -29,12 +31,20 @@
       @Property(description="Loop back multicast messages")
       boolean loopback=true;
   
  -    /* --------------------------------------------- Fields ------------------------------------------------------ */
  -    protected Address              local_addr, next;
  -    protected int                  view_size=0;
  -    protected Executor             default_pool=null;
  -    protected Executor             oob_pool=null;
  +    @Property
  +    int forward_queue_size=10000;
  +
  +    @Property
  +    int send_queue_size=10000;
   
  +    /* --------------------------------------------- Fields ------------------------------------------------------ */
  +    protected Address                local_addr, next;
  +    protected int                    view_size=0;
  +    protected Executor               default_pool=null;
  +    protected Executor               oob_pool=null;
  +    protected BlockingQueue<Message> send_queue;
  +    protected BlockingQueue<Message> forward_queue;
  +    protected boolean                forward=false; // flipped between true and false, to ensure fairness
   
       @ManagedAttribute
       public int msgs_forwarded=0;
  @@ -42,11 +52,17 @@
       @ManagedAttribute
       public int msgs_sent=0;
   
  +    @ManagedAttribute
  +    public int getForwardQueueSize() {return forward_queue.size();}
   
  +    @ManagedAttribute
  +    public int getSendQueueSize() {return send_queue.size();}
   
       public void init() throws Exception {
           default_pool=getTransport().getDefaultThreadPool();
           oob_pool=getTransport().getOOBThreadPool();
  +        send_queue=new ConcurrentLinkedBlockingQueue<Message>(send_queue_size);
  +        forward_queue=new ConcurrentLinkedBlockingQueue<Message>(forward_queue_size);
       }
   
       public Object down(final Event evt) {
  @@ -68,6 +84,15 @@
                   copy.setDest(next);
                   copy.putHeader(getId(), hdr);
   
  +                try {
  +                    msgs_sent++;
  +                    send_queue.put(copy);
  +                }
  +                catch(InterruptedException e) {
  +                    Thread.currentThread().interrupt();
  +                    return null;
  +                }
  +
                   if(loopback) {
                       if(log.isTraceEnabled()) log.trace(new StringBuilder("looping back message ").append(msg));
                       msg.setSrc(local_addr);
  @@ -80,8 +105,7 @@
                       });
                   }
   
  -                msgs_sent++;
  -                return forward(copy);
  +                return processQueues();
   
   
               case Event.VIEW_CHANGE:
  @@ -116,8 +140,13 @@
                       Message copy=msg.copy(true);
                       copy.setDest(next);
                       copy.putHeader(getId(), new DaisyHeader(ttl));
  -                    forward(copy);
  -                    msgs_forwarded++;
  +                    try {
  +                        msgs_forwarded++;
  +                        forward_queue.put(copy);
  +                    }
  +                    catch(InterruptedException e) {
  +                    }
  +                    processQueues();
                   }
   
                   // 2. Pass up
  @@ -128,17 +157,32 @@
       }
   
   
  -    protected Object forward(Message msg) {
  -        if(msg == null)
  -            return null;
  -
  -        if(log.isTraceEnabled()) {
  -            DaisyHeader hdr=(DaisyHeader)msg.getHeader(getId());
  -            log.trace(local_addr + ": forwarding message with ttl=" + hdr.getTTL() + " to " + next);
  +    protected Object processQueues() {
  +        for(int i=0; i < 10; i++) {
  +            try {
  +                Message msg=forward? forward_queue.poll() : send_queue.poll();
  +                if(msg == null) {
  +                    msg=forward? send_queue.poll() : forward_queue.poll();
  +                    if(msg == null)
  +                        continue;
  +                }
  +                if(log.isTraceEnabled()) {
  +                    DaisyHeader hdr=(DaisyHeader)msg.getHeader(getId());
  +                    log.trace(local_addr + ": " + (forward? " forwarding" : " sending") + " message with ttl=" + hdr.getTTL() + " to " + next);
  +                }
  +                return down_prot.down(new Event(Event.MSG, msg));
  +            }
  +            catch(Throwable t) {
  +                log.error("failed sending message down", t);
  +                return null;
  +            }
  +            finally {
  +                forward=!forward;
  +            }
           }
  -        return down_prot.down(new Event(Event.MSG, msg));
  +        return null;
       }
  -    
  +
   
       protected void handleView(View view) {
           view_size=view.size();
  
  
  

------------------------------------------------------------------------------
Start uncovering the many advantages of virtual appliances
and start using them to simplify application deployment and
accelerate your shift to cloud computing.
http://p.sf.net/sfu/novell-sfdev2dev