Class FanOutQueueImpl

  • All Implemented Interfaces:
    Closeable, AutoCloseable, IFanOutQueue

    public class FanOutQueueImpl
    extends Object
    implements IFanOutQueue
    A big, fast and persistent queue implementation supporting fan out semantics. 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. 6. FANOUT - support fan out semantics, multiple consumers can independently consume a single queue without intervention, everyone has its own queue front index. 7. CLIENT MANAGED INDEX - support access by index and the queue index is managed at client side.
    Author:
    bulldog
    • Method Summary

      All Methods Instance Methods Concrete Methods 
      Modifier and Type Method Description
      void close()  
      byte[] dequeue​(String fanoutId)
      Retrieves and removes the front of a fan out queue
      byte[] dequeue​(String fanoutId, boolean useLatest)
      Retrieves and removes the front of a fan out queue
      long enqueue​(byte[] data)
      Adds an item at the back of the queue
      long findClosestIndex​(long timestamp)
      Find an index closest to the specific timestamp when the corresponding item was enqueued.
      void flush()
      Force to persist current state of the queue, normally, you don't need to flush explicitly since: 1.)
      byte[] get​(long index)
      Retrieves data item at the specific index of the queue
      long getBackFileSize()
      Current total size of the back files of this queue
      long getFrontIndex()
      Get the queue front index, this is the earliest appended index
      long getFrontIndex​(String fanoutId)
      Get front index of specific fanout queue
      long getFrontIndex​(String fanoutId, boolean useLatest)
      Get front index of specific fanout queue
      int getLength​(long index)
      Get length of data item at specific index of the queue
      long getRearIndex()
      Get the queue rear index, this is the next to be appended index
      long getTimestamp​(long index)
      Get timestamp of data item at specific index of the queue, this is the timestamp when corresponding item was appended into the queue.
      boolean isEmpty()
      Determines whether the queue is empty
      boolean isEmpty​(String fanoutId)
      Determines whether a fan out queue is empty
      boolean isEmpty​(String fanoutId, boolean useLatest)
      Determines whether a fan out queue is empty
      void limitBackFileSize​(long sizeLimit)
      Limit the back file size of this queue, truncate back files and advance the queue front if necessary.
      byte[] peek​(String fanoutId)
      Peek the item at the front of a fanout queue, without removing it from the queue
      byte[] peek​(String fanoutId, boolean useLatest)
      Peek the item at the front of a fanout queue, without removing it from the queue
      int peekLength​(String fanoutId)
      Peek the length of the item at the front of a fan out queue
      int peekLength​(String fanoutId, boolean useLatest)
      Peek the length of the item at the front of a fan out queue
      long peekTimestamp​(String fanoutId)
      Peek the timestamp of the item at the front of a fan out queue
      long peekTimestamp​(String fanoutId, boolean useLatest)
      Peek the timestamp of the item at the front of a fan out queue
      void removeAll()
      Removes all items of a queue, this will empty the queue and delete all back data files.
      void removeBefore​(long timestamp)
      Remove all data before specific timestamp, truncate back files and advance the queue front if necessary.
      void resetQueueFrontIndex​(String fanoutId, long index)
      Reset the front index of a fanout queue.
      long size()
      Total number of items remaining in the queue.
      long size​(String fanoutId)
      Total number of items remaining in the fan out queue
      long size​(String fanoutId, boolean useLatest)
      Total number of items remaining in the fan out queue
    • Constructor Detail

      • FanOutQueueImpl

        public FanOutQueueImpl​(String queueDir,
                               String queueName,
                               int pageSize)
                        throws IOException
        A big, fast and persistent queue implementation with fandout support.
        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
      • FanOutQueueImpl

        public FanOutQueueImpl​(String queueDir,
                               String queueName)
                        throws IOException
        A big, fast and persistent queue implementation with fanout support, 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
    • Method Detail

      • isEmpty

        public boolean isEmpty​(String fanoutId,
                               boolean useLatest)
                        throws IOException
        Description copied from interface: IFanOutQueue
        Determines whether a fan out queue is empty
        Specified by:
        isEmpty in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        useLatest - if no offset has been recorded the head of the queue is used
        Returns:
        true if empty, false otherwise
        Throws:
        IOException - exception thrown if IO error occurs
      • isEmpty

        public boolean isEmpty​(String fanoutId)
                        throws IOException
        Description copied from interface: IFanOutQueue
        Determines whether a fan out queue is empty
        Specified by:
        isEmpty in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        Returns:
        true if empty, false otherwise
        Throws:
        IOException - exception thrown if IO error occurs
      • isEmpty

        public boolean isEmpty()
        Description copied from interface: IFanOutQueue
        Determines whether the queue is empty
        Specified by:
        isEmpty in interface IFanOutQueue
        Returns:
        true if empty, false otherwise
      • enqueue

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

        public byte[] dequeue​(String fanoutId)
                       throws IOException
        Description copied from interface: IFanOutQueue
        Retrieves and removes the front of a fan out queue
        Specified by:
        dequeue in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        Returns:
        data at the front of a queue
        Throws:
        IOException - exception throws if there is any IO error during dequeue operation.
      • dequeue

        public byte[] dequeue​(String fanoutId,
                              boolean useLatest)
                       throws IOException
        Description copied from interface: IFanOutQueue
        Retrieves and removes the front of a fan out queue
        Specified by:
        dequeue in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        useLatest - if no offset has been recorded the head of the queue is used
        Returns:
        data at the front of a queue
        Throws:
        IOException - exception throws if there is any IO error during dequeue operation.
      • peek

        public byte[] peek​(String fanoutId)
                    throws IOException
        Description copied from interface: IFanOutQueue
        Peek the item at the front of a fanout queue, without removing it from the queue
        Specified by:
        peek in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        Returns:
        data at the front of a queue
        Throws:
        IOException - exception throws if there is any IO error during peek operation.
      • peek

        public byte[] peek​(String fanoutId,
                           boolean useLatest)
                    throws IOException
        Description copied from interface: IFanOutQueue
        Peek the item at the front of a fanout queue, without removing it from the queue
        Specified by:
        peek in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        useLatest - if no offset has been recorded the head of the queue is used
        Returns:
        data at the front of a queue
        Throws:
        IOException - exception throws if there is any IO error during peek operation.
      • peekLength

        public int peekLength​(String fanoutId)
                       throws IOException
        Description copied from interface: IFanOutQueue
        Peek the length of the item at the front of a fan out queue
        Specified by:
        peekLength in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        Returns:
        data at the front of a queue
        Throws:
        IOException - exception throws if there is any IO error during peek operation.
      • peekLength

        public int peekLength​(String fanoutId,
                              boolean useLatest)
                       throws IOException
        Description copied from interface: IFanOutQueue
        Peek the length of the item at the front of a fan out queue
        Specified by:
        peekLength in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        useLatest - if no offset has been recorded the head of the queue is used
        Returns:
        data at the front of a queue
        Throws:
        IOException - exception throws if there is any IO error during peek operation.
      • peekTimestamp

        public long peekTimestamp​(String fanoutId)
                           throws IOException
        Description copied from interface: IFanOutQueue
        Peek the timestamp of the item at the front of a fan out queue
        Specified by:
        peekTimestamp in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        Returns:
        data at the front of a queue
        Throws:
        IOException - exception throws if there is any IO error during peek operation.
      • peekTimestamp

        public long peekTimestamp​(String fanoutId,
                                  boolean useLatest)
                           throws IOException
        Description copied from interface: IFanOutQueue
        Peek the timestamp of the item at the front of a fan out queue
        Specified by:
        peekTimestamp in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        useLatest - if no offset has been recorded the head of the queue is used
        Returns:
        data at the front of a queue
        Throws:
        IOException - exception throws if there is any IO error during peek operation.
      • get

        public byte[] get​(long index)
                   throws IOException
        Description copied from interface: IFanOutQueue
        Retrieves data item at the specific index of the queue
        Specified by:
        get in interface IFanOutQueue
        Parameters:
        index - data item index
        Returns:
        data at index
        Throws:
        IOException - exception throws if there is any IO error during fetch operation.
      • getLength

        public int getLength​(long index)
                      throws IOException
        Description copied from interface: IFanOutQueue
        Get length of data item at specific index of the queue
        Specified by:
        getLength in interface IFanOutQueue
        Parameters:
        index - data item index
        Returns:
        length of data item
        Throws:
        IOException - exception throws if there is any IO error during fetch operation.
      • getTimestamp

        public long getTimestamp​(long index)
                          throws IOException
        Description copied from interface: IFanOutQueue
        Get timestamp of data item at specific index of the queue, this is the timestamp when corresponding item was appended into the queue.
        Specified by:
        getTimestamp in interface IFanOutQueue
        Parameters:
        index - data item index
        Returns:
        timestamp of data item
        Throws:
        IOException - exception throws if there is any IO error during fetch operation.
      • size

        public long size​(String fanoutId)
                  throws IOException
        Description copied from interface: IFanOutQueue
        Total number of items remaining in the fan out queue
        Specified by:
        size in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        Returns:
        total number
        Throws:
        IOException - exception thrown if IO error occurs
      • removeBefore

        public void removeBefore​(long timestamp)
                          throws IOException
        Description copied from interface: IFanOutQueue
        Remove all data before specific timestamp, truncate back files and advance the queue front if necessary.
        Specified by:
        removeBefore in interface IFanOutQueue
        Parameters:
        timestamp - a timestamp
        Throws:
        IOException - exception thrown if there was any IO error during the removal operation
      • limitBackFileSize

        public void limitBackFileSize​(long sizeLimit)
                               throws IOException
        Description copied from interface: IFanOutQueue
        Limit the back file size of this queue, truncate back files and advance the queue front if necessary. Note, this is a best effort call, exact size limit can't be guaranteed
        Specified by:
        limitBackFileSize in interface IFanOutQueue
        Parameters:
        sizeLimit - size limit
        Throws:
        IOException - exception thrown if there was any IO error during the operation
      • getBackFileSize

        public long getBackFileSize()
                             throws IOException
        Description copied from interface: IFanOutQueue
        Current total size of the back files of this queue
        Specified by:
        getBackFileSize in interface IFanOutQueue
        Returns:
        total back file size
        Throws:
        IOException - exception thrown if there was any IO error during the operation
      • findClosestIndex

        public long findClosestIndex​(long timestamp)
                              throws IOException
        Description copied from interface: IFanOutQueue
        Find an index closest to the specific timestamp when the corresponding item was enqueued. to find latest index, use IFanOutQueue.LATEST as timestamp. to find earliest index, use IFanOutQueue.EARLIEST as timestamp.
        Specified by:
        findClosestIndex in interface IFanOutQueue
        Parameters:
        timestamp - when the corresponding item was appended
        Returns:
        an index
        Throws:
        IOException - exception thrown during the operation
      • resetQueueFrontIndex

        public void resetQueueFrontIndex​(String fanoutId,
                                         long index)
                                  throws IOException
        Description copied from interface: IFanOutQueue
        Reset the front index of a fanout queue.
        Specified by:
        resetQueueFrontIndex in interface IFanOutQueue
        Parameters:
        fanoutId - fanout identifier
        index - target index
        Throws:
        IOException - exception thrown during the operation
      • size

        public long size​(String fanoutId,
                         boolean useLatest)
                  throws IOException
        Description copied from interface: IFanOutQueue
        Total number of items remaining in the fan out queue
        Specified by:
        size in interface IFanOutQueue
        Parameters:
        fanoutId - the fanout identifier
        useLatest - if no offset has been recorded the head of the queue is used
        Returns:
        total number
        Throws:
        IOException - exception thrown if IO error occurs
      • size

        public long size()
        Description copied from interface: IFanOutQueue
        Total number of items remaining in the queue.
        Specified by:
        size in interface IFanOutQueue
        Returns:
        total number
      • flush

        public void flush()
        Description copied from interface: IFanOutQueue
        Force to persist current state of the queue, normally, you don't need to flush explicitly since: 1.) FanOutQueue will automatically flush a cached page when it is replaced out, 2.) FanOutQueue 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 IFanOutQueue
      • removeAll

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

        public long getFrontIndex()
        Description copied from interface: IFanOutQueue
        Get the queue front index, this is the earliest appended index
        Specified by:
        getFrontIndex in interface IFanOutQueue
        Returns:
        an index
      • getFrontIndex

        public long getFrontIndex​(String fanoutId)
                           throws IOException
        Description copied from interface: IFanOutQueue
        Get front index of specific fanout queue
        Specified by:
        getFrontIndex in interface IFanOutQueue
        Parameters:
        fanoutId - fanout identifier
        Returns:
        an index
        Throws:
        IOException - exception thrown if IO error occurs
      • getRearIndex

        public long getRearIndex()
        Description copied from interface: IFanOutQueue
        Get the queue rear index, this is the next to be appended index
        Specified by:
        getRearIndex in interface IFanOutQueue
        Returns:
        an index
      • getFrontIndex

        public long getFrontIndex​(String fanoutId,
                                  boolean useLatest)
                           throws IOException
        Description copied from interface: IFanOutQueue
        Get front index of specific fanout queue
        Specified by:
        getFrontIndex in interface IFanOutQueue
        Parameters:
        fanoutId - fanout identifier
        useLatest - if no offset has been recorded the head of the queue is used
        Returns:
        an index
        Throws:
        IOException - exception thrown if IO error occurs