public class FanOutQueue extends Object implements Closeable
| Modifier and Type | Field and Description |
|---|---|
static long |
EARLIEST
Constant represents earliest timestamp
|
static long |
LATEST
Constant represents latest timestamp
|
| Constructor and Description |
|---|
FanOutQueue(String queueDir,
String queueName)
A big, fast and persistent queue implementation with fanout support,
use default back data page size, see
BigArray.DEFAULT_DATA_PAGE_SIZE |
FanOutQueue(String queueDir,
String queueName,
int pageSize)
A big, fast and persistent queue implementation with fandout support.
|
| Modifier and Type | Method and Description |
|---|---|
void |
close() |
byte[] |
dequeue(String fanoutId)
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.) 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.
|
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) |
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) |
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
|
int |
peekLength(String fanoutId)
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
|
void |
removeAll() |
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) |
long |
size()
Total number of items remaining in the queue.
|
long |
size(String fanoutId) |
public static final long EARLIEST
public static final long LATEST
public FanOutQueue(String queueDir, String queueName, int pageSize)
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 allowed BigArray.MINIMUM_DATA_PAGE_SIZEpublic FanOutQueue(String queueDir, String queueName)
BigArray.DEFAULT_DATA_PAGE_SIZEqueueDir - the directory to store queue dataqueueName - the name of the queue, will be appended as last part of the queue directorypublic boolean isEmpty(String fanoutId)
public boolean isEmpty()
public long enqueue(byte[] data)
data - to be enqueued datapublic byte[] dequeue(String fanoutId)
fanoutId - the fanout identifierpublic byte[] peek(String fanoutId)
fanoutId - the fanout identifierpublic int peekLength(String fanoutId)
fanoutId - the fanout identifierpublic long peekTimestamp(String fanoutId)
fanoutId - the fanout identifierpublic byte[] get(long index)
index - data item indexpublic int getLength(long index)
index - data item indexpublic long getTimestamp(long index)
index - data item indexpublic void removeBefore(long timestamp)
timestamp - a timestamppublic void limitBackFileSize(long sizeLimit)
sizeLimit - size limitpublic long getBackFileSize()
public long findClosestIndex(long timestamp)
LATEST as timestamp.
to find earliest index, use EARLIEST as timestamp.timestamp - when the corresponding item was appendedpublic void resetQueueFrontIndex(String fanoutId, long index)
public long size(String fanoutId)
public long size()
public void flush()
public void close()
close in interface Closeableclose in interface AutoCloseablepublic void removeAll()
public long getFrontIndex()
public long getRearIndex()
public long getFrontIndex(String fanoutId)
Copyright © 2016. All rights reserved.