Package org.kairosdb.bigqueue
Class BigQueueImpl
- java.lang.Object
-
- org.kairosdb.bigqueue.BigQueueImpl
-
- All Implemented Interfaces:
Closeable,AutoCloseable,IBigQueue
public class BigQueueImpl extends Object implements IBigQueue
A big, fast and persistent queue implementation. Main features: 1. FAST : close to the speed of direct memory access, both enqueue and dequeue are close to O(1) memory access. 2. MEMORY-EFFICIENT : automatic paging and swapping algorithm, only most-recently accessed data is kept in memory. 3. THREAD-SAFE : multiple threads can concurrently enqueue and dequeue without data corruption. 4. PERSISTENT - all data in queue is persisted on disk, and is crash resistant. 5. BIG(HUGE) - the total size of the queued data is only limited by the available disk space.- Author:
- bulldog
-
-
Nested Class Summary
-
Nested classes/interfaces inherited from interface org.kairosdb.bigqueue.IBigQueue
IBigQueue.ItemIterator
-
-
Constructor Summary
Constructors Constructor Description BigQueueImpl(String queueDir, String queueName)A big, fast and persistent queue implementation, use default back data page size, seeBigArrayImpl.DEFAULT_DATA_PAGE_SIZEBigQueueImpl(String queueDir, String queueName, int pageSize)A big, fast and persistent queue implementation.
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description voidapplyForEach(IBigQueue.ItemIterator iterator)apply an implementation of a ItemIterator interface for each queue itemvoidclose()byte[]dequeue()Retrieves and removes the front of a queueCompletableFuture<byte[]>dequeueAsync()Retrieves a Future which will complete if new Items where enqued.voidenqueue(byte[] data)Adds an item at the back of a queuevoidflush()Force to persist current state of the queue, normally, you don't need to flush explicitly since: 1.)voidgc()Delete all used data files to free disk space.booleanisEmpty()Determines whether a queue is emptybyte[]peek()Retrieves the item at the front of a queueCompletableFuture<byte[]>peekAsync()Retrieves the item at the front of a queue asynchronously.voidremoveAll()Removes all items of a queue, this will empty the queue and delete all back data files.longsize()Total number of items available in the queue.
-
-
-
Constructor Detail
-
BigQueueImpl
public BigQueueImpl(String queueDir, String queueName) throws IOException
A big, fast and persistent queue implementation, use default back data page size, seeBigArrayImpl.DEFAULT_DATA_PAGE_SIZE- Parameters:
queueDir- the directory to store queue dataqueueName- the name of the queue, will be appended as last part of the queue directory- Throws:
IOException- exception throws if there is any IO error during queue initialization
-
BigQueueImpl
public BigQueueImpl(String queueDir, String queueName, int pageSize) throws IOException
A big, fast and persistent queue implementation.- Parameters:
queueDir- the directory to store queue dataqueueName- the name of the queue, will be appended as last part of the queue directorypageSize- the back data file size per page in bytes, see minimum allowedBigArrayImpl.MINIMUM_DATA_PAGE_SIZE- Throws:
IOException- exception throws if there is any IO error during queue initialization
-
-
Method Detail
-
isEmpty
public boolean isEmpty()
Description copied from interface:IBigQueueDetermines whether a queue is empty
-
enqueue
public void enqueue(byte[] data) throws IOExceptionDescription copied from interface:IBigQueueAdds an item at the back of a queue- Specified by:
enqueuein interfaceIBigQueue- Parameters:
data- to be enqueued data- Throws:
IOException- exception throws if there is any IO error during enqueue operation.
-
dequeue
public byte[] dequeue() throws IOExceptionDescription copied from interface:IBigQueueRetrieves and removes the front of a queue- Specified by:
dequeuein interfaceIBigQueue- Returns:
- data at the front of a queue
- Throws:
IOException- exception throws if there is any IO error during dequeue operation.
-
dequeueAsync
public CompletableFuture<byte[]> dequeueAsync()
Description copied from interface:IBigQueueRetrieves a Future which will complete if new Items where enqued. Use this method to retrieve a future where to register as Listener instead of repeatedly polling the queues state. On complete this future contains the result of the dequeue operation. Hence the item was automatically removed from the queue.- Specified by:
dequeueAsyncin interfaceIBigQueue- Returns:
- a ListenableFuture which completes with the first entry if items are ready to be dequeued.
-
removeAll
public void removeAll() throws IOExceptionDescription copied from interface:IBigQueueRemoves all items of a queue, this will empty the queue and delete all back data files.- Specified by:
removeAllin interfaceIBigQueue- Throws:
IOException- exception throws if there is any IO error during dequeue operation.
-
peek
public byte[] peek() throws IOExceptionDescription copied from interface:IBigQueueRetrieves the item at the front of a queue- Specified by:
peekin interfaceIBigQueue- Returns:
- data at the front of a queue
- Throws:
IOException- exception throws if there is any IO error during peek operation.
-
peekAsync
public CompletableFuture<byte[]> peekAsync()
Description copied from interface:IBigQueueRetrieves the item at the front of a queue asynchronously. On complete the value set in this future is the result of the peek operation. Hence the item remains at the front of the list.
-
applyForEach
public void applyForEach(IBigQueue.ItemIterator iterator) throws IOException
apply an implementation of a ItemIterator interface for each queue item- Specified by:
applyForEachin interfaceIBigQueue- Parameters:
iterator- Callback used for each item in array.- Throws:
IOException- exception thrown if IO error occurs
-
close
public void close() throws IOException- Specified by:
closein interfaceAutoCloseable- Specified by:
closein interfaceCloseable- Throws:
IOException
-
gc
public void gc() throws IOExceptionDescription copied from interface:IBigQueueDelete all used data files to free disk space. BigQueue will persist enqueued data in disk files, these data files will remain even after the data in them has been dequeued later, so your application is responsible to periodically call this method to delete all used data files and free disk space.- Specified by:
gcin interfaceIBigQueue- Throws:
IOException- exception throws if there is any IO error during gc operation.
-
flush
public void flush()
Description copied from interface:IBigQueueForce to persist current state of the queue, normally, you don't need to flush explicitly since: 1.) BigQueue will automatically flush a cached page when it is replaced out, 2.) BigQueue uses memory mapped file technology internally, and the OS will flush the changes even your process crashes, call this periodically only if you need transactional reliability and you are aware of the cost to performance.
-
-