public abstract class AbstractTableReaderStream extends AbstractStream implements Runnable, DbReaderStream
| Modifier and Type | Field and Description |
|---|---|
protected boolean |
ascending |
protected boolean |
follow |
protected PartitionManager |
partitionManager |
protected boolean |
quit |
protected TableDefinition |
tableDefinition |
log, name, outputDefinition, state, subscribers, ydb| Modifier | Constructor and Description |
|---|---|
protected |
AbstractTableReaderStream(YarchDatabase ydb,
TableDefinition tblDef,
PartitionManager partitionManager,
boolean ascending,
boolean follow) |
| Modifier and Type | Method and Description |
|---|---|
boolean |
addInFilter(ColumnExpression cexpr,
Set<Object> values)
currently adds only filters on value based partitions
|
boolean |
addRelOpFilter(ColumnExpression cexpr,
RelOp relOp,
Object value) |
protected int |
compare(byte[] a1,
byte[] a2) |
protected Tuple |
dataToTuple(byte[] k,
byte[] v) |
void |
doClose() |
protected boolean |
emitIfNotPastStart(byte[] key,
byte[] value,
byte[] rangeStart,
boolean strictStart) |
protected boolean |
emitIfNotPastStop(byte[] key,
byte[] value,
byte[] rangeEnd,
boolean strictEnd) |
TableDefinition |
getTableDefinition() |
void |
run() |
protected abstract boolean |
runPartitions(List<Partition> partitions,
IndexFilter range)
Runs the partitions sending data only that conform with the start and end filters.
|
void |
start()
Start emitting tuples.
|
addSubscriber, close, emitTuple, getColumnDefinition, getDefinition, getName, getNumEmittedTuples, getState, getSubscriberCount, getSubscribers, removeSubscriber, setName, toStringprotected TableDefinition tableDefinition
protected volatile boolean quit
protected final PartitionManager partitionManager
protected final boolean ascending
protected final boolean follow
protected AbstractTableReaderStream(YarchDatabase ydb, TableDefinition tblDef, PartitionManager partitionManager, boolean ascending, boolean follow)
public void start()
Streamstart in interface Streamstart in class AbstractStreamprotected abstract boolean runPartitions(List<Partition> partitions, IndexFilter range) throws IOException
IOExceptionprotected boolean emitIfNotPastStop(byte[] key,
byte[] value,
byte[] rangeEnd,
boolean strictEnd)
protected boolean emitIfNotPastStart(byte[] key,
byte[] value,
byte[] rangeStart,
boolean strictStart)
public boolean addRelOpFilter(ColumnExpression cexpr, RelOp relOp, Object value) throws StreamSqlException
addRelOpFilter in interface DbReaderStreamStreamSqlExceptionprotected Tuple dataToTuple(byte[] k, byte[] v)
public boolean addInFilter(ColumnExpression cexpr, Set<Object> values) throws StreamSqlException
addInFilter in interface DbReaderStreamStreamSqlExceptionpublic void doClose()
doClose in class AbstractStreampublic TableDefinition getTableDefinition()
protected int compare(byte[] a1,
byte[] a2)
Copyright © 2017 Space Applications Services. All rights reserved.