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]>