Package org.kairosdb.bigqueue
Class FanOutQueueImpl
- java.lang.Object
-
- org.kairosdb.bigqueue.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
-
-
Field Summary
-
Fields inherited from interface org.kairosdb.bigqueue.IFanOutQueue
EARLIEST, LATEST
-
-
Constructor Summary
Constructors Constructor Description FanOutQueueImpl(String queueDir, String queueName)A big, fast and persistent queue implementation with fanout support, use default back data page size, seeBigArrayImpl.DEFAULT_DATA_PAGE_SIZEFanOutQueueImpl(String queueDir, String queueName, int pageSize)A big, fast and persistent queue implementation with fandout support.
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description voidclose()byte[]dequeue(String fanoutId)Retrieves and removes the front of a fan out queuebyte[]dequeue(String fanoutId, boolean useLatest)Retrieves and removes the front of a fan out queuelongenqueue(byte[] data)Adds an item at the back of the queuelongfindClosestIndex(long timestamp)Find an index closest to the specific timestamp when the corresponding item was enqueued.voidflush()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 queuelonggetBackFileSize()Current total size of the back files of this queuelonggetFrontIndex()Get the queue front index, this is the earliest appended indexlonggetFrontIndex(String fanoutId)Get front index of specific fanout queuelonggetFrontIndex(String fanoutId, boolean useLatest)Get front index of specific fanout queueintgetLength(long index)Get length of data item at specific index of the queuelonggetRearIndex()Get the queue rear index, this is the next to be appended indexlonggetTimestamp(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.booleanisEmpty()Determines whether the queue is emptybooleanisEmpty(String fanoutId)Determines whether a fan out queue is emptybooleanisEmpty(String fanoutId, boolean useLatest)Determines whether a fan out queue is emptyvoidlimitBackFileSize(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 queuebyte[]peek(String fanoutId, boolean useLatest)Peek the item at the front of a fanout queue, without removing it from the queueintpeekLength(String fanoutId)Peek the length of the item at the front of a fan out queueintpeekLength(String fanoutId, boolean useLatest)Peek the length of the item at the front of a fan out queuelongpeekTimestamp(String fanoutId)Peek the timestamp of the item at the front of a fan out queuelongpeekTimestamp(String fanoutId, boolean useLatest)Peek the timestamp of the item at the front of a fan out queuevoidremoveAll()Removes all items of a queue, this will empty the queue and delete all back data files.voidremoveBefore(long timestamp)Remove all data before specific timestamp, truncate back files and advance the queue front if necessary.voidresetQueueFrontIndex(String fanoutId, long index)Reset the front index of a fanout queue.longsize()Total number of items remaining in the queue.longsize(String fanoutId)Total number of items remaining in the fan out queuelongsize(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 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
-
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, 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
-
-
Method Detail
-
isEmpty
public boolean isEmpty(String fanoutId, boolean useLatest) throws IOException
Description copied from interface:IFanOutQueueDetermines whether a fan out queue is empty- Specified by:
isEmptyin interfaceIFanOutQueue- Parameters:
fanoutId- the fanout identifieruseLatest- 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:IFanOutQueueDetermines whether a fan out queue is empty- Specified by:
isEmptyin interfaceIFanOutQueue- 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:IFanOutQueueDetermines whether the queue is empty- Specified by:
isEmptyin interfaceIFanOutQueue- Returns:
- true if empty, false otherwise
-
enqueue
public long enqueue(byte[] data) throws IOExceptionDescription copied from interface:IFanOutQueueAdds an item at the back of the queue- Specified by:
enqueuein interfaceIFanOutQueue- 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:IFanOutQueueRetrieves and removes the front of a fan out queue- Specified by:
dequeuein interfaceIFanOutQueue- 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:IFanOutQueueRetrieves and removes the front of a fan out queue- Specified by:
dequeuein interfaceIFanOutQueue- Parameters:
fanoutId- the fanout identifieruseLatest- 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:IFanOutQueuePeek the item at the front of a fanout queue, without removing it from the queue- Specified by:
peekin interfaceIFanOutQueue- 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:IFanOutQueuePeek the item at the front of a fanout queue, without removing it from the queue- Specified by:
peekin interfaceIFanOutQueue- Parameters:
fanoutId- the fanout identifieruseLatest- 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:IFanOutQueuePeek the length of the item at the front of a fan out queue- Specified by:
peekLengthin interfaceIFanOutQueue- 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:IFanOutQueuePeek the length of the item at the front of a fan out queue- Specified by:
peekLengthin interfaceIFanOutQueue- Parameters:
fanoutId- the fanout identifieruseLatest- 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:IFanOutQueuePeek the timestamp of the item at the front of a fan out queue- Specified by:
peekTimestampin interfaceIFanOutQueue- 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:IFanOutQueuePeek the timestamp of the item at the front of a fan out queue- Specified by:
peekTimestampin interfaceIFanOutQueue- Parameters:
fanoutId- the fanout identifieruseLatest- 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 IOExceptionDescription copied from interface:IFanOutQueueRetrieves data item at the specific index of the queue- Specified by:
getin interfaceIFanOutQueue- 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 IOExceptionDescription copied from interface:IFanOutQueueGet length of data item at specific index of the queue- Specified by:
getLengthin interfaceIFanOutQueue- 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 IOExceptionDescription copied from interface:IFanOutQueueGet timestamp of data item at specific index of the queue, this is the timestamp when corresponding item was appended into the queue.- Specified by:
getTimestampin interfaceIFanOutQueue- 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:IFanOutQueueTotal number of items remaining in the fan out queue- Specified by:
sizein interfaceIFanOutQueue- Parameters:
fanoutId- the fanout identifier- Returns:
- total number
- Throws:
IOException- exception thrown if IO error occurs
-
removeBefore
public void removeBefore(long timestamp) throws IOExceptionDescription copied from interface:IFanOutQueueRemove all data before specific timestamp, truncate back files and advance the queue front if necessary.- Specified by:
removeBeforein interfaceIFanOutQueue- 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 IOExceptionDescription copied from interface:IFanOutQueueLimit 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:
limitBackFileSizein interfaceIFanOutQueue- Parameters:
sizeLimit- size limit- Throws:
IOException- exception thrown if there was any IO error during the operation
-
getBackFileSize
public long getBackFileSize() throws IOExceptionDescription copied from interface:IFanOutQueueCurrent total size of the back files of this queue- Specified by:
getBackFileSizein interfaceIFanOutQueue- 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 IOExceptionDescription copied from interface:IFanOutQueueFind an index closest to the specific timestamp when the corresponding item was enqueued. to find latest index, useIFanOutQueue.LATESTas timestamp. to find earliest index, useIFanOutQueue.EARLIESTas timestamp.- Specified by:
findClosestIndexin interfaceIFanOutQueue- 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:IFanOutQueueReset the front index of a fanout queue.- Specified by:
resetQueueFrontIndexin interfaceIFanOutQueue- Parameters:
fanoutId- fanout identifierindex- target index- Throws:
IOException- exception thrown during the operation
-
size
public long size(String fanoutId, boolean useLatest) throws IOException
Description copied from interface:IFanOutQueueTotal number of items remaining in the fan out queue- Specified by:
sizein interfaceIFanOutQueue- Parameters:
fanoutId- the fanout identifieruseLatest- 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:IFanOutQueueTotal number of items remaining in the queue.- Specified by:
sizein interfaceIFanOutQueue- Returns:
- total number
-
flush
public void flush()
Description copied from interface:IFanOutQueueForce 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:
flushin interfaceIFanOutQueue
-
close
public void close() throws IOException- Specified by:
closein interfaceAutoCloseable- Specified by:
closein interfaceCloseable- Throws:
IOException
-
removeAll
public void removeAll() throws IOExceptionDescription copied from interface:IFanOutQueueRemoves all items of a queue, this will empty the queue and delete all back data files.- Specified by:
removeAllin interfaceIFanOutQueue- Throws:
IOException- exception throws if there is any IO error during dequeue operation.
-
getFrontIndex
public long getFrontIndex()
Description copied from interface:IFanOutQueueGet the queue front index, this is the earliest appended index- Specified by:
getFrontIndexin interfaceIFanOutQueue- Returns:
- an index
-
getFrontIndex
public long getFrontIndex(String fanoutId) throws IOException
Description copied from interface:IFanOutQueueGet front index of specific fanout queue- Specified by:
getFrontIndexin interfaceIFanOutQueue- Parameters:
fanoutId- fanout identifier- Returns:
- an index
- Throws:
IOException- exception thrown if IO error occurs
-
getRearIndex
public long getRearIndex()
Description copied from interface:IFanOutQueueGet the queue rear index, this is the next to be appended index- Specified by:
getRearIndexin interfaceIFanOutQueue- Returns:
- an index
-
getFrontIndex
public long getFrontIndex(String fanoutId, boolean useLatest) throws IOException
Description copied from interface:IFanOutQueueGet front index of specific fanout queue- Specified by:
getFrontIndexin interfaceIFanOutQueue- Parameters:
fanoutId- fanout identifieruseLatest- if no offset has been recorded the head of the queue is used- Returns:
- an index
- Throws:
IOException- exception thrown if IO error occurs
-
-