public class RecordOutputCollector extends Object implements org.apache.flink.util.Collector<org.apache.flink.types.Record>
Records, and emits the pair to a set of Nephele RecordWriters.
The OutputCollector tracks to which writers a deep-copy must be given and which not.| Modifier and Type | Field and Description |
|---|---|
protected RecordWriter<org.apache.flink.types.Record>[] |
writers |
| Constructor and Description |
|---|
RecordOutputCollector(List<RecordWriter<org.apache.flink.types.Record>> writers)
Initializes the output collector with a set of writers.
|
| Modifier and Type | Method and Description |
|---|---|
void |
addWriter(RecordWriter<org.apache.flink.types.Record> writer)
Adds a writer to the OutputCollector.
|
void |
close() |
void |
collect(org.apache.flink.types.Record record)
Collects a
Record, and emits it to all writers. |
List<RecordWriter<org.apache.flink.types.Record>> |
getWriters()
List of writers that are associated with this output collector
|
protected RecordWriter<org.apache.flink.types.Record>[] writers
public RecordOutputCollector(List<RecordWriter<org.apache.flink.types.Record>> writers)
List.writers - List of all writers.public void addWriter(RecordWriter<org.apache.flink.types.Record> writer)
writer - The writer to add.public void collect(org.apache.flink.types.Record record)
Record, and emits it to all writers.
Writers which require a deep-copy are fed with a copy.collect in interface org.apache.flink.util.Collector<org.apache.flink.types.Record>public void close()
close in interface org.apache.flink.util.Collector<org.apache.flink.types.Record>public List<RecordWriter<org.apache.flink.types.Record>> getWriters()
Copyright © 2014–2015 The Apache Software Foundation. All rights reserved.