public class AccumulatorManager extends Object
| Constructor and Description |
|---|
AccumulatorManager(int maxEntries) |
| Modifier and Type | Method and Description |
|---|---|
Map<String,Object> |
getJobAccumulatorResults(org.apache.flink.api.common.JobID jobID) |
Map<String,SerializedValue<Object>> |
getJobAccumulatorResultsSerialized(org.apache.flink.api.common.JobID jobID) |
StringifiedAccumulatorResult[] |
getJobAccumulatorResultsStringified(org.apache.flink.api.common.JobID jobID) |
void |
processIncomingAccumulators(org.apache.flink.api.common.JobID jobID,
Map<String,org.apache.flink.api.common.accumulators.Accumulator<?,?>> newAccumulators)
Merges the new accumulators with the existing accumulators collected for
the job.
|
public void processIncomingAccumulators(org.apache.flink.api.common.JobID jobID,
Map<String,org.apache.flink.api.common.accumulators.Accumulator<?,?>> newAccumulators)
public Map<String,Object> getJobAccumulatorResults(org.apache.flink.api.common.JobID jobID)
public Map<String,SerializedValue<Object>> getJobAccumulatorResultsSerialized(org.apache.flink.api.common.JobID jobID) throws IOException
IOExceptionpublic StringifiedAccumulatorResult[] getJobAccumulatorResultsStringified(org.apache.flink.api.common.JobID jobID) throws IOException
IOExceptionCopyright © 2014–2015 The Apache Software Foundation. All rights reserved.