public class RealtimePlumber extends Object implements Plumber
Constructor and Description |
---|
RealtimePlumber(DataSchema schema,
RealtimeTuningConfig config,
FireDepartmentMetrics metrics,
com.metamx.emitter.service.ServiceEmitter emitter,
QueryRunnerFactoryConglomerate conglomerate,
DataSegmentAnnouncer segmentAnnouncer,
ExecutorService queryExecutorService,
io.druid.segment.loading.DataSegmentPusher dataSegmentPusher,
SegmentPublisher segmentPublisher,
FilteredServerView serverView) |
Modifier and Type | Method and Description |
---|---|
protected void |
abandonSegment(long truncatedTime,
Sink sink)
Unannounces a given sink and removes all local references to it.
|
int |
add(io.druid.data.input.InputRow row) |
protected void |
bootstrapSinksFromDisk() |
protected File |
computeBaseDir(DataSchema schema) |
protected File |
computePersistDir(DataSchema schema,
org.joda.time.Interval interval) |
void |
finishJob()
Perform any final processing and clean up after ourselves.
|
RealtimeTuningConfig |
getConfig() |
<T> QueryRunner<T> |
getQueryRunner(Query<T> query) |
RejectionPolicy |
getRejectionPolicy() |
DataSchema |
getSchema() |
Sink |
getSink(long timestamp) |
Map<Long,Sink> |
getSinks() |
protected void |
initializeExecutors() |
void |
persist(Runnable commitRunnable)
Persist any in-memory indexed data to durable storage.
|
protected int |
persistHydrant(FireHydrant indexToPersist,
DataSchema schema,
org.joda.time.Interval interval)
Persists the given hydrant and returns the number of rows persisted
|
protected void |
shutdownExecutors() |
void |
startJob()
Perform any initial setup.
|
protected void |
startPersistThread() |
public RealtimePlumber(DataSchema schema, RealtimeTuningConfig config, FireDepartmentMetrics metrics, com.metamx.emitter.service.ServiceEmitter emitter, QueryRunnerFactoryConglomerate conglomerate, DataSegmentAnnouncer segmentAnnouncer, ExecutorService queryExecutorService, io.druid.segment.loading.DataSegmentPusher dataSegmentPusher, SegmentPublisher segmentPublisher, FilteredServerView serverView)
public DataSchema getSchema()
public RealtimeTuningConfig getConfig()
public RejectionPolicy getRejectionPolicy()
public void startJob()
Plumber
Plumber.finishJob()
.public int add(io.druid.data.input.InputRow row) throws IndexSizeExceededException
add
in interface Plumber
row
- - the row to insertIndexSizeExceededException
public <T> QueryRunner<T> getQueryRunner(Query<T> query)
getQueryRunner
in interface Plumber
public void persist(Runnable commitRunnable)
Plumber
public void finishJob()
Plumber
protected void initializeExecutors()
protected void shutdownExecutors()
protected void bootstrapSinksFromDisk()
protected void startPersistThread()
protected void abandonSegment(long truncatedTime, Sink sink)
truncatedTime
- sink keysink
- sink to unannounceprotected File computeBaseDir(DataSchema schema)
protected File computePersistDir(DataSchema schema, org.joda.time.Interval interval)
protected int persistHydrant(FireHydrant indexToPersist, DataSchema schema, org.joda.time.Interval interval)
indexToPersist
- hydrant to persistschema
- datasource schemainterval
- interval to persistCopyright © 2011–2015. All rights reserved.