public abstract class AbstractStream extends Object implements Stream
| Modifier and Type | Field and Description |
|---|---|
protected org.slf4j.Logger |
log |
protected String |
name |
protected TupleDefinition |
outputDefinition |
protected int |
state |
protected Collection<StreamSubscriber> |
subscribers |
protected YarchDatabase |
ydb |
| Modifier | Constructor and Description |
|---|---|
protected |
AbstractStream(YarchDatabase ydb,
String name,
TupleDefinition definition) |
| Modifier and Type | Method and Description |
|---|---|
void |
addSubscriber(StreamSubscriber s) |
void |
close()
Closes the stream by:
send the streamClosed signal to all subscribed clients
|
protected abstract void |
doClose() |
void |
emitTuple(Tuple t) |
ColumnDefinition |
getColumnDefinition(String colName) |
TupleDefinition |
getDefinition() |
String |
getName() |
long |
getNumEmittedTuples() |
int |
getState() |
int |
getSubscriberCount() |
Collection<StreamSubscriber> |
getSubscribers() |
void |
removeSubscriber(StreamSubscriber s) |
void |
setName(String streamName) |
abstract void |
start()
Start emitting tuples.
|
String |
toString() |
protected String name
protected TupleDefinition outputDefinition
protected final Collection<StreamSubscriber> subscribers
protected volatile int state
protected org.slf4j.Logger log
protected YarchDatabase ydb
protected AbstractStream(YarchDatabase ydb, String name, TupleDefinition definition)
public abstract void start()
Streampublic TupleDefinition getDefinition()
getDefinition in interface Streampublic void addSubscriber(StreamSubscriber s)
addSubscriber in interface Streampublic void removeSubscriber(StreamSubscriber s)
removeSubscriber in interface Streampublic ColumnDefinition getColumnDefinition(String colName)
getColumnDefinition in interface Streampublic final void close()
Streamprotected abstract void doClose()
public long getNumEmittedTuples()
getNumEmittedTuples in interface Streampublic int getSubscriberCount()
getSubscriberCount in interface Streampublic Collection<StreamSubscriber> getSubscribers()
getSubscribers in interface StreamCopyright © 2017 Space Applications Services. All rights reserved.