bloritsch    01/12/18 13:09:28

  Modified:    src/scratchpad/org/apache/avalon/excalibur/event
                        DefaultQueue.java
               src/scratchpad/org/apache/avalon/excalibur/event/test
                        QueueTestCase.java
  Log:
  threadsafe queue implementation can perform over 11 million enqueue/dequeue 
ops in 12.7 secs
  
  Revision  Changes    Path
  1.2       +93 -21    
jakarta-avalon-excalibur/src/scratchpad/org/apache/avalon/excalibur/event/DefaultQueue.java
  
  Index: DefaultQueue.java
  ===================================================================
  RCS file: 
/home/cvs/jakarta-avalon-excalibur/src/scratchpad/org/apache/avalon/excalibur/event/DefaultQueue.java,v
  retrieving revision 1.1
  retrieving revision 1.2
  diff -u -r1.1 -r1.2
  --- DefaultQueue.java 2001/12/17 23:04:59     1.1
  +++ DefaultQueue.java 2001/12/18 21:09:28     1.2
  @@ -8,19 +8,23 @@
   package org.apache.avalon.excalibur.event;
   
   import java.util.ArrayList;
  +import org.apache.avalon.excalibur.concurrent.Mutex;
   
   /**
  - * The default queue implementation is a variabl size queue.
  + * The default queue implementation is a variabl size queue.  This queue is
  + * ThreadSafe, however the overhead in synchronization costs a few extra 
millis.
    *
    * @author <a href="mailto:[EMAIL PROTECTED]">Berin Loritsch</a>
    */
  -public class DefaultQueue extends AbstractQueue
  +public final class DefaultQueue extends AbstractQueue
   {
       private final ArrayList m_elements;
  +    private final Mutex     m_mutex;
   
       public DefaultQueue()
       {
           m_elements = new ArrayList();
  +        m_mutex = new Mutex();
       }
   
       public int size()
  @@ -36,8 +40,19 @@
   
       public boolean tryEnqueue( final QueueElement element )
       {
  -        boolean success = m_elements.add( element );
  -        m_elements.notifyAll();
  +        boolean success = false;
  +        try
  +        {
  +            m_mutex.acquire();
  +            success = m_elements.add( element );
  +        }
  +        catch ( InterruptedException ie )
  +        {
  +        }
  +        finally
  +        {
  +            m_mutex.release();
  +        }
   
           return success;
       }
  @@ -47,24 +62,43 @@
       {
           final int len = elements.length;
   
  -        for ( int i = 0; i < len; i++ )
  +        try
           {
  -            m_elements.add( elements[i] );
  -        }
  +            m_mutex.acquire();
   
  -        m_elements.notifyAll();
  +            for ( int i = 0; i < len; i++ )
  +            {
  +                m_elements.add( elements[i] );
  +            }
  +        }
  +        catch ( InterruptedException ie )
  +        {
  +        }
  +        finally
  +        {
  +            m_mutex.release();
  +        }
       }
   
       public void enqueue( final QueueElement element )
           throws SourceException
       {
  -        m_elements.add( element );
  +        try
  +        {
  +            m_mutex.acquire();
  +            m_elements.add( element );
  +        }
  +        catch ( InterruptedException ie )
  +        {
  +        }
  +        finally
  +        {
  +            m_mutex.release();
  +        }
       }
   
       public QueueElement[] dequeue( final int numElements )
       {
  -        block( m_elements );
  -
           int arraySize = numElements;
   
           if ( size() < numElements )
  @@ -72,36 +106,74 @@
               arraySize = size();
           }
   
  -        QueueElement[] elements = new QueueElement[ arraySize ];
  +        QueueElement[] elements = null;
   
  -        for ( int i = 0; i < arraySize; i++ )
  +        try
           {
  -            elements[i] = (QueueElement) m_elements.remove( 0 );
  +            m_mutex.attempt( m_timeout );
  +
  +            elements = new QueueElement[ arraySize ];
  +
  +            for ( int i = 0; i < arraySize; i++ )
  +            {
  +                elements[i] = (QueueElement) m_elements.remove( 0 );
  +            }
           }
  +        catch ( InterruptedException ie )
  +        {
  +        }
  +        finally
  +        {
  +            m_mutex.release();
  +        }
   
           return elements;
       }
   
       public QueueElement[] dequeueAll()
       {
  -        block( m_elements );
  +        QueueElement[] elements = null;
  +
  +        try
  +        {
  +            m_mutex.attempt( m_timeout );
   
  -        QueueElement[] elements = (QueueElement[]) m_elements.toArray( new 
QueueElement [] {} );
  -        m_elements.clear();
  +            elements = (QueueElement[]) m_elements.toArray( new QueueElement 
[] {} );
  +            m_elements.clear();
  +        }
  +        catch ( InterruptedException ie )
  +        {
  +        }
  +        finally
  +        {
  +            m_mutex.release();
  +        }
   
           return elements;
       }
   
       public QueueElement dequeue()
       {
  -        block( m_elements );
  +        QueueElement element = null;
   
  -        if ( size() <= 0 )
  +        try
  +        {
  +            m_mutex.attempt( m_timeout );
  +
  +            if ( size() > 0 )
  +            {
  +                element = (QueueElement) m_elements.remove( 0 );
  +            }
  +        }
  +        catch ( InterruptedException ie )
  +        {
  +        }
  +        finally
           {
  -            return null;
  +            m_mutex.release();
           }
   
  -        return (QueueElement) m_elements.remove( 0 );
  +        return element;
       }
   
       private final static class DefaultPreparedEnqueue implements 
PreparedEnqueue
  
  
  
  1.2       +68 -1     
jakarta-avalon-excalibur/src/scratchpad/org/apache/avalon/excalibur/event/test/QueueTestCase.java
  
  Index: QueueTestCase.java
  ===================================================================
  RCS file: 
/home/cvs/jakarta-avalon-excalibur/src/scratchpad/org/apache/avalon/excalibur/event/test/QueueTestCase.java,v
  retrieving revision 1.1
  retrieving revision 1.2
  diff -u -r1.1 -r1.2
  --- QueueTestCase.java        2001/12/17 23:05:00     1.1
  +++ QueueTestCase.java        2001/12/18 21:09:28     1.2
  @@ -9,6 +9,7 @@
   
   import junit.framework.TestCase;
   import org.apache.avalon.excalibur.event.DefaultQueue;
  +import org.apache.avalon.excalibur.event.PreparedEnqueue;
   import org.apache.avalon.excalibur.event.QueueElement;
   import org.apache.avalon.excalibur.event.SourceException;
   
  @@ -19,6 +20,9 @@
    */
   public class QueueTestCase extends TestCase
   {
  +    QueueElement element = new TestQueueElement();
  +    QueueElement[] elements = new TestQueueElement[10];
  +
       private static final class TestQueueElement implements QueueElement
       {
            private final Object m_attachment;
  @@ -47,8 +51,46 @@
       public QueueTestCase( String name )
       {
           super( name );
  +
  +        for ( int i = 0; i < 10; i++ )
  +        {
  +            elements[i] = new TestQueueElement();
  +        }
  +    }
  +
  +    public void testMillionIterationOneElement()
  +        throws Exception
  +    {
  +        DefaultQueue queue = new DefaultQueue();
  +        assertEquals( queue.size(), 0 );
  +
  +        for ( int j = 0; j < 1000000; j++ )
  +        {
  +            queue.enqueue( element );
  +            assertEquals( queue.size(), 1 );
  +
  +            assertNotNull( queue.dequeue() );
  +            assertEquals( queue.size(), 0 );
  +        }
       }
   
  +    public void testMillionIterationTenElements()
  +        throws Exception
  +    {
  +        DefaultQueue queue = new DefaultQueue();
  +        assertEquals( queue.size(), 0 );
  +
  +        for ( int j = 0; j < 1000000; j++ )
  +        {
  +            queue.enqueue( elements );
  +            assertEquals( queue.size(), 10 );
  +
  +            QueueElement[] results = queue.dequeueAll();
  +            assertEquals( results.length, 10 );
  +            assertEquals( queue.size(), 0 );
  +        }
  +    }
  +
       public void testDefaultQueue()
       {
           DefaultQueue queue = new DefaultQueue();
  @@ -58,9 +100,34 @@
           try
           {
               queue.enqueue( new TestQueueElement () );
  -            assertTrue( queue.size() > 0 );
  +            assertEquals( queue.size(), 1 );
   
               assertNotNull( queue.dequeue() );
  +            assertEquals( queue.size(), 0 );
  +
  +            queue.enqueue( elements );
  +            assertEquals( queue.size(), 10 );
  +
  +            QueueElement[] results = queue.dequeue( 3 );
  +            assertEquals( results.length, 3 );
  +            assertEquals( queue.size(), 7 );
  +
  +            results = queue.dequeueAll();
  +            assertEquals( results.length, 7 );
  +            assertEquals( queue.size(), 0 );
  +
  +            PreparedEnqueue prep = queue.prepareEnqueue( elements );
  +            assertEquals( queue.size(), 0 );
  +            prep.abort();
  +            assertEquals( queue.size(), 0 );
  +
  +            prep = queue.prepareEnqueue( elements );
  +            assertEquals( queue.size(), 0 );
  +            prep.commit();
  +            assertEquals( queue.size(), 10 );
  +
  +            results = queue.dequeue( queue.size() );
  +            assertEquals( queue.size(), 0 );
           }
           catch ( SourceException se )
           {
  
  
  

--
To unsubscribe, e-mail:   <mailto:[EMAIL PROTECTED]>
For additional commands, e-mail: <mailto:[EMAIL PROTECTED]>

Reply via email to