public class GraphStreamQuery extends GraphQuery implements io.reactivex.ObservableOnSubscribe<commonj.sdo.DataGraph>, org.plasma.sdo.access.StreamQueryDispatcher
context| Constructor and Description |
|---|
GraphStreamQuery(ServiceContext context) |
| Modifier and Type | Method and Description |
|---|---|
void |
close() |
protected ResultsAssembler |
createResultsAssembler(org.plasma.query.model.Query query,
org.plasma.query.collector.SelectionCollector selection,
Expr whereSyntaxTree,
ResultsComparator orderingComparator,
ResultsComparator groupingComparator,
Expr havingSyntaxTree,
TableReader rootTableReader,
GraphAssemblerFactory assemblerFactory) |
protected void |
execute(org.apache.hadoop.hbase.client.Scan scan,
TableReader rootTableReader,
ResultsAssembler collector) |
protected void |
executeAsStream(org.plasma.query.model.Query query,
org.plasma.query.collector.SelectionCollector selection,
org.plasma.sdo.PlasmaType type,
org.apache.hadoop.hbase.filter.Filter columnFilter,
Expr whereSyntaxTree,
ResultsComparator orderingComparator,
List<PartialRowKey> partialScans,
List<FuzzyRowKey> fuzzyScans,
List<CompleteRowKey> completeKeys,
Timestamp snapshotDate) |
io.reactivex.Observable<commonj.sdo.DataGraph> |
findAsStream(org.plasma.query.model.Query query,
int requestMax,
Timestamp snapshotDate) |
io.reactivex.Observable<commonj.sdo.DataGraph> |
findAsStream(org.plasma.query.model.Query query,
Timestamp snapshotDate) |
void |
subscribe(io.reactivex.ObservableEmitter<commonj.sdo.DataGraph> emitter) |
canAbortScan, collectRowKeyProperties, count, createRootColumnFilterAssembler, createScan, createScan, execute, execute, execute, execute, execute, execute, execute, find, find, getVariables, log, serializeGraphpublic GraphStreamQuery(ServiceContext context)
public void close()
close in interface org.plasma.sdo.access.QueryDispatcherclose in interface org.plasma.sdo.access.StreamQueryDispatcherclose in class GraphQuerypublic io.reactivex.Observable<commonj.sdo.DataGraph> findAsStream(org.plasma.query.model.Query query,
Timestamp snapshotDate)
findAsStream in interface org.plasma.sdo.access.StreamQueryDispatcherpublic io.reactivex.Observable<commonj.sdo.DataGraph> findAsStream(org.plasma.query.model.Query query,
int requestMax,
Timestamp snapshotDate)
findAsStream in interface org.plasma.sdo.access.StreamQueryDispatcherpublic void subscribe(io.reactivex.ObservableEmitter<commonj.sdo.DataGraph> emitter)
throws Exception
subscribe in interface io.reactivex.ObservableOnSubscribe<commonj.sdo.DataGraph>Exceptionprotected void executeAsStream(org.plasma.query.model.Query query,
org.plasma.query.collector.SelectionCollector selection,
org.plasma.sdo.PlasmaType type,
org.apache.hadoop.hbase.filter.Filter columnFilter,
Expr whereSyntaxTree,
ResultsComparator orderingComparator,
List<PartialRowKey> partialScans,
List<FuzzyRowKey> fuzzyScans,
List<CompleteRowKey> completeKeys,
Timestamp snapshotDate)
protected void execute(org.apache.hadoop.hbase.client.Scan scan,
TableReader rootTableReader,
ResultsAssembler collector)
throws IOException
execute in class GraphQueryIOExceptionprotected ResultsAssembler createResultsAssembler(org.plasma.query.model.Query query, org.plasma.query.collector.SelectionCollector selection, Expr whereSyntaxTree, ResultsComparator orderingComparator, ResultsComparator groupingComparator, Expr havingSyntaxTree, TableReader rootTableReader, GraphAssemblerFactory assemblerFactory)
createResultsAssembler in class GraphQueryCopyright © 2021. All Rights Reserved.