Class 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
    • Constructor Detail

      • BigQueueImpl

        public BigQueueImpl​(String queueDir,
                            String queueName)
                     throws IOException
        A big, fast and persistent queue implementation, use default back data page size, see BigArrayImpl.DEFAULT_DATA_PAGE_SIZE
        Parameters:
        queueDir - the directory to store queue data
        queueName - 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 data
        queueName - the name of the queue, will be appended as last part of the queue directory
        pageSize - the back data file size per page in bytes, see minimum allowed BigArrayImpl.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: IBigQueue
        Determines whether a queue is empty
        Specified by:
        isEmpty in interface IBigQueue
        Returns:
        ture if empty, false otherwise
      • enqueue

        public void enqueue​(byte[] data)
                     throws IOException
        Description copied from interface: IBigQueue
        Adds an item at the back of a queue
        Specified by:
        enqueue in interface IBigQueue
        Parameters:
        data - to be enqueued data
        Throws:
        IOException - exception throws if there is any IO error during enqueue operation.
      • dequeue

        public byte[] dequeue()
                       throws IOException
        Description copied from interface: IBigQueue
        Retrieves and removes the front of a queue
        Specified by:
        dequeue in interface IBigQueue
        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: IBigQueue
        Retrieves 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:
        dequeueAsync in interface IBigQueue
        Returns:
        a ListenableFuture which completes with the first entry if items are ready to be dequeued.
      • removeAll

        public void removeAll()
                       throws IOException
        Description copied from interface: IBigQueue
        Removes all items of a queue, this will empty the queue and delete all back data files.
        Specified by:
        removeAll in interface IBigQueue
        Throws:
        IOException - exception throws if there is any IO error during dequeue operation.
      • peek

        public byte[] peek()
                    throws IOException
        Description copied from interface: IBigQueue
        Retrieves the item at the front of a queue
        Specified by:
        peek in interface IBigQueue
        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: IBigQueue
        Retrieves 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.
        Specified by:
        peekAsync in interface IBigQueue
        Returns:
        a future containing the first item if available. You may register as listener at this future to be informed if a new item arrives.
      • applyForEach

        public void applyForEach​(IBigQueue.ItemIterator iterator)
                          throws IOException
        apply an implementation of a ItemIterator interface for each queue item
        Specified by:
        applyForEach in interface IBigQueue
        Parameters:
        iterator - Callback used for each item in array.
        Throws:
        IOException - exception thrown if IO error occurs
      • gc

        public void gc()
                throws IOException
        Description copied from interface: IBigQueue
        Delete 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:
        gc in interface IBigQueue
        Throws:
        IOException - exception throws if there is any IO error during gc operation.
      • flush

        public void flush()
        Description copied from interface: IBigQueue
        Force 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.
        Specified by:
        flush in interface IBigQueue
      • size

        public long size()
        Description copied from interface: IBigQueue
        Total number of items available in the queue.
        Specified by:
        size in interface IBigQueue
        Returns:
        total number