public abstract class AbstractTableReaderStream extends AbstractStream implements Runnable, DbReaderStream
Stream.ExceptionHandler| 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(YarchDatabaseInstance 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.
|
addSubscriber, close, emitTuple, exceptionHandler, getColumnDefinition, getDefinition, getName, getNumEmittedTuples, getState, getSubscriberCount, getSubscribers, removeSubscriber, setName, start, toStringprotected TableDefinition tableDefinition
protected volatile boolean quit
protected final PartitionManager partitionManager
protected final boolean ascending
protected final boolean follow
protected AbstractTableReaderStream(YarchDatabaseInstance ydb, TableDefinition tblDef, PartitionManager partitionManager, boolean ascending, boolean follow)
protected 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 © 2018 Space Applications Services. All rights reserved.