Modifier and Type | Method and Description |
---|---|
<K,V> DataSource<Tuple2<K,V>> |
ExecutionEnvironment.createHadoopInput(org.apache.hadoop.mapreduce.InputFormat<K,V> mapreduceInputFormat,
Class<K> key,
Class<V> value,
org.apache.hadoop.mapreduce.Job job)
Creates a
DataSet from the given InputFormat . |
<K,V> DataSource<Tuple2<K,V>> |
ExecutionEnvironment.createHadoopInput(org.apache.hadoop.mapred.InputFormat<K,V> mapredInputFormat,
Class<K> key,
Class<V> value,
org.apache.hadoop.mapred.JobConf job)
Creates a
DataSet from the given InputFormat . |
<K,V> DataSource<Tuple2<K,V>> |
ExecutionEnvironment.readHadoopFile(org.apache.hadoop.mapred.FileInputFormat<K,V> mapredInputFormat,
Class<K> key,
Class<V> value,
String inputPath)
Creates a
DataSet from the given FileInputFormat . |
<K,V> DataSource<Tuple2<K,V>> |
ExecutionEnvironment.readHadoopFile(org.apache.hadoop.mapreduce.lib.input.FileInputFormat<K,V> mapreduceInputFormat,
Class<K> key,
Class<V> value,
String inputPath)
Creates a
DataSet from the given FileInputFormat . |
<K,V> DataSource<Tuple2<K,V>> |
ExecutionEnvironment.readHadoopFile(org.apache.hadoop.mapreduce.lib.input.FileInputFormat<K,V> mapreduceInputFormat,
Class<K> key,
Class<V> value,
String inputPath,
org.apache.hadoop.mapreduce.Job job)
Creates a
DataSet from the given FileInputFormat . |
<K,V> DataSource<Tuple2<K,V>> |
ExecutionEnvironment.readHadoopFile(org.apache.hadoop.mapred.FileInputFormat<K,V> mapredInputFormat,
Class<K> key,
Class<V> value,
String inputPath,
org.apache.hadoop.mapred.JobConf job)
Creates a
DataSet from the given FileInputFormat . |
<K,V> DataSource<Tuple2<K,V>> |
ExecutionEnvironment.readSequenceFile(Class<K> key,
Class<V> value,
String inputPath)
|
Modifier and Type | Method and Description |
---|---|
Tuple2<K,V> |
HadoopInputFormat.nextRecord(Tuple2<K,V> record) |
Modifier and Type | Method and Description |
---|---|
TypeInformation<Tuple2<K,V>> |
HadoopInputFormat.getProducedType() |
Modifier and Type | Method and Description |
---|---|
Tuple2<K,V> |
HadoopInputFormat.nextRecord(Tuple2<K,V> record) |
void |
HadoopOutputFormat.writeRecord(Tuple2<K,V> record) |
Modifier and Type | Method and Description |
---|---|
Tuple2<K,V> |
HadoopInputFormat.nextRecord(Tuple2<K,V> record) |
Modifier and Type | Method and Description |
---|---|
TypeInformation<Tuple2<K,V>> |
HadoopInputFormat.getProducedType() |
Modifier and Type | Method and Description |
---|---|
Tuple2<K,V> |
HadoopInputFormat.nextRecord(Tuple2<K,V> record) |
void |
HadoopOutputFormat.writeRecord(Tuple2<K,V> record) |
Modifier and Type | Method and Description |
---|---|
<T0,T1> DataSource<Tuple2<T0,T1>> |
CsvReader.types(Class<T0> type0,
Class<T1> type1)
Specifies the types for the CSV fields.
|
Modifier and Type | Method and Description |
---|---|
Tuple2<T1,T2> |
CrossOperator.DefaultCrossFunction.cross(T1 first,
T2 second) |
Modifier and Type | Method and Description |
---|---|
static <T,K> Operator<Tuple2<K,T>> |
KeyFunctions.appendKeyExtractor(Operator<T> input,
Keys.SelectorFunctionKeys<T,K> key) |
static <T,K> TypeInformation<Tuple2<K,T>> |
KeyFunctions.createTypeWithKey(Keys.SelectorFunctionKeys<T,K> key) |
<T0,T1> ProjectOperator<T,Tuple2<T0,T1>> |
ProjectOperator.Projection.projectTuple2()
|
<T0,T1> JoinOperator.ProjectJoin<I1,I2,Tuple2<T0,T1>> |
JoinOperator.JoinProjection.projectTuple2()
Projects a pair of joined elements to a
Tuple with the previously selected fields. |
<T0,T1> CrossOperator.ProjectCross<I1,I2,Tuple2<T0,T1>> |
CrossOperator.CrossProjection.projectTuple2()
Projects a pair of crossed elements to a
Tuple with the previously selected fields. |
Modifier and Type | Method and Description |
---|---|
static <T,K> SingleInputOperator<?,T,?> |
KeyFunctions.appendKeyRemover(Operator<Tuple2<K,T>> inputWithKey,
Keys.SelectorFunctionKeys<T,K> key) |
void |
JoinOperator.DefaultFlatJoinFunction.join(T1 first,
T2 second,
Collector<Tuple2<T1,T2>> out) |
Modifier and Type | Method and Description |
---|---|
Tuple2<K,T> |
KeyExtractingMapper.map(T value) |
Tuple2<K,T> |
PlanUnwrappingReduceOperator.ReduceWrapper.reduce(Tuple2<K,T> value1,
Tuple2<K,T> value2) |
Modifier and Type | Method and Description |
---|---|
void |
TupleRightUnwrappingJoiner.join(I1 value1,
Tuple2<K,I2> value2,
Collector<OUT> collector) |
void |
TupleLeftUnwrappingJoiner.join(Tuple2<K,I1> value1,
I2 value2,
Collector<OUT> collector) |
void |
TupleUnwrappingJoiner.join(Tuple2<K,I1> value1,
Tuple2<K,I2> value2,
Collector<OUT> collector) |
void |
TupleUnwrappingJoiner.join(Tuple2<K,I1> value1,
Tuple2<K,I2> value2,
Collector<OUT> collector) |
T |
KeyRemovingMapper.map(Tuple2<K,T> value) |
Tuple2<K,T> |
PlanUnwrappingReduceOperator.ReduceWrapper.reduce(Tuple2<K,T> value1,
Tuple2<K,T> value2) |
Tuple2<K,T> |
PlanUnwrappingReduceOperator.ReduceWrapper.reduce(Tuple2<K,T> value1,
Tuple2<K,T> value2) |
Modifier and Type | Method and Description |
---|---|
void |
PlanRightUnwrappingCoGroupOperator.TupleRightUnwrappingCoGrouper.coGroup(Iterable<I1> records1,
Iterable<Tuple2<K,I2>> records2,
Collector<OUT> out) |
void |
PlanLeftUnwrappingCoGroupOperator.TupleLeftUnwrappingCoGrouper.coGroup(Iterable<Tuple2<K,I1>> records1,
Iterable<I2> records2,
Collector<OUT> out) |
void |
PlanBothUnwrappingCoGroupOperator.TupleBothUnwrappingCoGrouper.coGroup(Iterable<Tuple2<K,I1>> records1,
Iterable<Tuple2<K,I2>> records2,
Collector<OUT> out) |
void |
PlanBothUnwrappingCoGroupOperator.TupleBothUnwrappingCoGrouper.coGroup(Iterable<Tuple2<K,I1>> records1,
Iterable<Tuple2<K,I2>> records2,
Collector<OUT> out) |
void |
PlanUnwrappingGroupCombineOperator.TupleUnwrappingGroupCombiner.combine(Iterable<Tuple2<K,IN>> values,
Collector<OUT> out) |
void |
PlanUnwrappingReduceGroupOperator.TupleUnwrappingGroupCombinableGroupReducer.combine(Iterable<Tuple2<K,IN>> values,
Collector<Tuple2<K,IN>> out) |
void |
PlanUnwrappingReduceGroupOperator.TupleUnwrappingGroupCombinableGroupReducer.combine(Iterable<Tuple2<K,IN>> values,
Collector<Tuple2<K,IN>> out) |
void |
PlanUnwrappingReduceGroupOperator.TupleUnwrappingGroupCombinableGroupReducer.reduce(Iterable<Tuple2<K,IN>> values,
Collector<OUT> out) |
void |
PlanUnwrappingReduceGroupOperator.TupleUnwrappingNonCombinableGroupReducer.reduce(Iterable<Tuple2<K,IN>> values,
Collector<OUT> out) |
void |
TupleWrappingCollector.set(Collector<Tuple2<K,IN>> wrappedCollector) |
void |
TupleUnwrappingIterator.set(Iterator<Tuple2<K,T>> iterator) |
Modifier and Type | Method and Description |
---|---|
Tuple2<T0,T1> |
Tuple2.copy()
Shallow tuple copy.
|
static <T0,T1> Tuple2<T0,T1> |
Tuple2.of(T0 value0,
T1 value1)
Creates a new tuple and assigns the given values to the tuple's fields.
|
Tuple2<T1,T0> |
Tuple2.swap()
Returns a shallow copy of the tuple with swapped values.
|
Modifier and Type | Method and Description |
---|---|
Tuple2<T0,T1>[] |
Tuple2Builder.build() |
Modifier and Type | Method and Description |
---|---|
static <T> DataSet<Tuple2<Integer,Long>> |
DataSetUtils.countElementsPerPartition(DataSet<T> input)
Method that goes over all the elements in each partition in order to retrieve
the total number of elements.
|
static <T> DataSet<Tuple2<Long,T>> |
DataSetUtils.zipWithIndex(DataSet<T> input)
Method that assigns a unique
Long value to all elements in the input data set. |
static <T> DataSet<Tuple2<Long,T>> |
DataSetUtils.zipWithUniqueId(DataSet<T> input)
Method that assigns a unique
Long value to all elements in the input data set in the following way:
a map function is applied to the input data set
each map task holds a counter c which is increased for each record
c is shifted by n bits where n = log2(number of parallel tasks)
to create a unique ID among all tasks, the task id is added to the counter
for each record, the resulting counter is collected
|
Modifier and Type | Method and Description |
---|---|
Map<Tuple2<K,N>,com.google.common.base.Optional<V>> |
LazyDbValueState.getModified()
Return the Map of modified states that hasn't been written to the
database yet.
|
Map<Tuple2<K,N>,com.google.common.base.Optional<V>> |
LazyDbValueState.getStateCache()
Return the Map of cached states.
|
Modifier and Type | Method and Description |
---|---|
void |
DbAdapter.insertBatch(String stateId,
DbBackendConfig conf,
Connection con,
PreparedStatement insertStatement,
long checkpointTimestamp,
List<Tuple2<byte[],byte[]>> toInsert)
Insert a list of Key-Value pairs into the database.
|
void |
MySqlAdapter.insertBatch(String stateId,
DbBackendConfig conf,
Connection con,
PreparedStatement insertStatement,
long checkpointTs,
List<Tuple2<byte[],byte[]>> toInsert) |
Modifier and Type | Method and Description |
---|---|
Tuple2<Integer,KMeans.Point> |
KMeans.SelectNearestCenter.map(KMeans.Point p) |
Modifier and Type | Method and Description |
---|---|
Tuple3<Integer,KMeans.Point,Long> |
KMeans.CountAppender.map(Tuple2<Integer,KMeans.Point> t) |
Modifier and Type | Method and Description |
---|---|
Tuple2<Long,Long> |
ConnectedComponents.NeighborWithComponentIDJoin.join(Tuple2<Long,Long> vertexWithComponent,
Tuple2<Long,Long> edge) |
Tuple2<Long,Double> |
PageRank.RankAssigner.map(Long page) |
Tuple2<T,T> |
ConnectedComponents.DuplicateValue.map(T vertex) |
Tuple2<Long,Double> |
PageRank.Dampener.map(Tuple2<Long,Double> value) |
Modifier and Type | Method and Description |
---|---|
boolean |
PageRank.EpsilonFilter.filter(Tuple2<Tuple2<Long,Double>,Tuple2<Long,Double>> value) |
void |
ConnectedComponents.UndirectEdge.flatMap(Tuple2<Long,Long> edge,
Collector<Tuple2<Long,Long>> out) |
void |
PageRank.JoinVertexWithEdgesMatch.flatMap(Tuple2<Tuple2<Long,Double>,Tuple2<Long,Long[]>> value,
Collector<Tuple2<Long,Double>> out) |
Tuple2<Long,Long> |
ConnectedComponents.NeighborWithComponentIDJoin.join(Tuple2<Long,Long> vertexWithComponent,
Tuple2<Long,Long> edge) |
Tuple2<Long,Long> |
ConnectedComponents.NeighborWithComponentIDJoin.join(Tuple2<Long,Long> vertexWithComponent,
Tuple2<Long,Long> edge) |
void |
ConnectedComponents.ComponentIdFilter.join(Tuple2<Long,Long> candidate,
Tuple2<Long,Long> old,
Collector<Tuple2<Long,Long>> out) |
void |
ConnectedComponents.ComponentIdFilter.join(Tuple2<Long,Long> candidate,
Tuple2<Long,Long> old,
Collector<Tuple2<Long,Long>> out) |
EnumTrianglesDataTypes.Edge |
EnumTriangles.TupleEdgeConverter.map(Tuple2<Integer,Integer> t) |
Tuple2<Long,Double> |
PageRank.Dampener.map(Tuple2<Long,Double> value) |
Modifier and Type | Method and Description |
---|---|
boolean |
PageRank.EpsilonFilter.filter(Tuple2<Tuple2<Long,Double>,Tuple2<Long,Double>> value) |
boolean |
PageRank.EpsilonFilter.filter(Tuple2<Tuple2<Long,Double>,Tuple2<Long,Double>> value) |
void |
ConnectedComponents.UndirectEdge.flatMap(Tuple2<Long,Long> edge,
Collector<Tuple2<Long,Long>> out) |
void |
PageRank.JoinVertexWithEdgesMatch.flatMap(Tuple2<Tuple2<Long,Double>,Tuple2<Long,Long[]>> value,
Collector<Tuple2<Long,Double>> out) |
void |
PageRank.JoinVertexWithEdgesMatch.flatMap(Tuple2<Tuple2<Long,Double>,Tuple2<Long,Long[]>> value,
Collector<Tuple2<Long,Double>> out) |
void |
PageRank.JoinVertexWithEdgesMatch.flatMap(Tuple2<Tuple2<Long,Double>,Tuple2<Long,Long[]>> value,
Collector<Tuple2<Long,Double>> out) |
void |
ConnectedComponents.ComponentIdFilter.join(Tuple2<Long,Long> candidate,
Tuple2<Long,Long> old,
Collector<Tuple2<Long,Long>> out) |
void |
PageRank.BuildOutgoingEdgeList.reduce(Iterable<Tuple2<Long,Long>> values,
Collector<Tuple2<Long,Long[]>> out) |
void |
PageRank.BuildOutgoingEdgeList.reduce(Iterable<Tuple2<Long,Long>> values,
Collector<Tuple2<Long,Long[]>> out) |
Modifier and Type | Class and Description |
---|---|
static class |
EnumTrianglesDataTypes.Edge |
Modifier and Type | Method and Description |
---|---|
static DataSet<Tuple2<Long,Long>> |
PageRankData.getDefaultEdgeDataSet(ExecutionEnvironment env) |
static DataSet<Tuple2<Long,Long>> |
ConnectedComponentsData.getDefaultEdgeDataSet(ExecutionEnvironment env) |
Modifier and Type | Method and Description |
---|---|
void |
EnumTrianglesDataTypes.Edge.copyVerticesFromTuple2(Tuple2<Integer,Integer> t) |
Modifier and Type | Method and Description |
---|---|
Tuple2<LinearRegression.Params,Integer> |
LinearRegression.SubUpdate.map(LinearRegression.Data in) |
Tuple2<LinearRegression.Params,Integer> |
LinearRegression.UpdateAccumulator.reduce(Tuple2<LinearRegression.Params,Integer> val1,
Tuple2<LinearRegression.Params,Integer> val2) |
Modifier and Type | Method and Description |
---|---|
LinearRegression.Params |
LinearRegression.Update.map(Tuple2<LinearRegression.Params,Integer> arg0) |
Tuple2<LinearRegression.Params,Integer> |
LinearRegression.UpdateAccumulator.reduce(Tuple2<LinearRegression.Params,Integer> val1,
Tuple2<LinearRegression.Params,Integer> val2) |
Tuple2<LinearRegression.Params,Integer> |
LinearRegression.UpdateAccumulator.reduce(Tuple2<LinearRegression.Params,Integer> val1,
Tuple2<LinearRegression.Params,Integer> val2) |
Modifier and Type | Class and Description |
---|---|
static class |
TPCHQuery3.Customer |
Modifier and Type | Method and Description |
---|---|
boolean |
WebLogAnalysis.FilterDocByKeyWords.filter(Tuple2<String,String> value)
Filters for documents that contain all of the given keywords and projects the records on the URL field.
|
boolean |
WebLogAnalysis.FilterVisitsByDate.filter(Tuple2<String,String> value)
Filters for records of the visits relation where the year of visit is equal to a
specified value.
|
Modifier and Type | Method and Description |
---|---|
static DataSet<Tuple2<String,String>> |
WebLogData.getDocumentDataSet(ExecutionEnvironment env) |
static DataSet<Tuple2<String,String>> |
WebLogData.getVisitDataSet(ExecutionEnvironment env) |
Modifier and Type | Method and Description |
---|---|
void |
WordCount.Tokenizer.flatMap(String value,
Collector<Tuple2<String,Integer>> out) |
Modifier and Type | Class and Description |
---|---|
class |
Vertex<K,V>
Represents the graph's nodes.
|
Modifier and Type | Method and Description |
---|---|
DataSet<Tuple2<K,Long>> |
Graph.getDegrees()
Return the degree of all vertices in the graph
|
DataSet<Tuple2<K,K>> |
Graph.getEdgeIds() |
DataSet<Tuple2<K,VV>> |
Graph.getVerticesAsTuple2() |
DataSet<Tuple2<K,Long>> |
Graph.inDegrees()
Return the in-degree of all vertices in the graph
|
DataSet<Tuple2<K,Long>> |
Graph.outDegrees()
Return the out-degree of all vertices in the graph
|
DataSet<Tuple2<K,EV>> |
Graph.reduceOnEdges(ReduceEdgesFunction<EV> reduceEdgesFunction,
EdgeDirection direction)
Compute a reduce transformation over the edge values of each vertex.
|
DataSet<Tuple2<K,VV>> |
Graph.reduceOnNeighbors(ReduceNeighborsFunction<VV> reduceNeighborsFunction,
EdgeDirection direction)
Compute a reduce transformation over the neighbors' vertex values of each vertex.
|
Modifier and Type | Method and Description |
---|---|
static <K> Graph<K,NullValue,NullValue> |
Graph.fromTuple2DataSet(DataSet<Tuple2<K,K>> edges,
ExecutionEnvironment context)
Creates a graph from a DataSet of Tuple2 objects for edges.
|
static <K,VV> Graph<K,VV,NullValue> |
Graph.fromTuple2DataSet(DataSet<Tuple2<K,K>> edges,
MapFunction<K,VV> vertexValueInitializer,
ExecutionEnvironment context)
Creates a graph from a DataSet of Tuple2 objects for edges.
|
static <K,VV,EV> Graph<K,VV,EV> |
Graph.fromTupleDataSet(DataSet<Tuple2<K,VV>> vertices,
DataSet<Tuple3<K,K,EV>> edges,
ExecutionEnvironment context)
Creates a graph from a DataSet of Tuple2 objects for vertices and
Tuple3 objects for edges.
|
void |
EdgesFunction.iterateEdges(Iterable<Tuple2<K,Edge<K,EV>>> edges,
Collector<O> out)
This method is called per vertex and can iterate over all of its neighboring edges
with the specified direction.
|
void |
NeighborsFunctionWithVertexValue.iterateNeighbors(Vertex<K,VV> vertex,
Iterable<Tuple2<Edge<K,EV>,Vertex<K,VV>>> neighbors,
Collector<O> out)
This method is called per vertex and can iterate over all of its neighbors
with the specified direction.
|
<T> Graph<K,VV,EV> |
Graph.joinWithEdgesOnSource(DataSet<Tuple2<K,T>> inputDataSet,
EdgeJoinFunction<EV,T> edgeJoinFunction)
Joins the edge DataSet with an input Tuple2 DataSet and applies a user-defined transformation
on the values of the matched records.
|
<T> Graph<K,VV,EV> |
Graph.joinWithEdgesOnTarget(DataSet<Tuple2<K,T>> inputDataSet,
EdgeJoinFunction<EV,T> edgeJoinFunction)
Joins the edge DataSet with an input Tuple2 DataSet and applies a user-defined transformation
on the values of the matched records.
|
<T> Graph<K,VV,EV> |
Graph.joinWithVertices(DataSet<Tuple2<K,T>> inputDataSet,
VertexJoinFunction<VV,T> vertexJoinFunction)
Joins the vertex DataSet of this graph with an input Tuple2 DataSet and applies
a user-defined transformation on the values of the matched records.
|
Modifier and Type | Method and Description |
---|---|
boolean |
MusicProfiles.FilterSongNodes.filter(Tuple2<String,String> value) |
Modifier and Type | Method and Description |
---|---|
void |
MusicProfiles.GetTopSongPerUser.iterateEdges(Vertex<String,NullValue> vertex,
Iterable<Edge<String,Integer>> edges,
Collector<Tuple2<String,String>> out) |
Modifier and Type | Class and Description |
---|---|
class |
Neighbor<VV,EV>
This class represents a
<sourceVertex, edge> pair
This is a wrapper around Tuple2<VV, EV> for convenience in the GatherFunction |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<String,DataSet<?>>> |
GSAConfiguration.getApplyBcastVars()
Get the broadcast variables of the ApplyFunction.
|
List<Tuple2<String,DataSet<?>>> |
GSAConfiguration.getGatherBcastVars()
Get the broadcast variables of the GatherFunction.
|
List<Tuple2<String,DataSet<?>>> |
GSAConfiguration.getSumBcastVars()
Get the broadcast variables of the SumFunction.
|
Modifier and Type | Class and Description |
---|---|
static class |
Summarization.EdgeValue<EV>
Value that is stored at a summarized edge.
|
static class |
Summarization.VertexValue<VV>
Value that is stored at a summarized vertex.
|
static class |
Summarization.VertexWithRepresentative<K>
Represents a vertex identifier and its corresponding vertex group identifier.
|
Modifier and Type | Method and Description |
---|---|
Vertex<K,Tuple2<Long,Double>> |
CommunityDetection.AddScoreToVertexValuesMapper.map(Vertex<K,Long> vertex) |
Modifier and Type | Method and Description |
---|---|
Long |
CommunityDetection.RemoveScoreFromVertexValuesMapper.map(Vertex<K,Tuple2<Long,Double>> vertex) |
void |
CommunityDetection.LabelMessenger.sendMessages(Vertex<K,Tuple2<Long,Double>> vertex) |
void |
CommunityDetection.VertexLabelUpdater.updateVertex(Vertex<K,Tuple2<Long,Double>> vertex,
MessageIterator<Tuple2<Long,Double>> inMessages) |
void |
CommunityDetection.VertexLabelUpdater.updateVertex(Vertex<K,Tuple2<Long,Double>> vertex,
MessageIterator<Tuple2<Long,Double>> inMessages) |
Modifier and Type | Method and Description |
---|---|
void |
EdgesFunction.iterateEdges(Iterable<Tuple2<K,Edge<K,EV>>> edges,
Collector<T> out) |
void |
NeighborsFunctionWithVertexValue.iterateNeighbors(Vertex<K,VV> vertex,
Iterable<Tuple2<Edge<K,EV>,Vertex<K,VV>>> neighbors,
Collector<T> out) |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<String,DataSet<?>>> |
ScatterGatherConfiguration.getMessagingBcastVars()
Get the broadcast variables of the MessagingFunction.
|
List<Tuple2<String,DataSet<?>>> |
ScatterGatherConfiguration.getUpdateBcastVars()
Get the broadcast variables of the VertexUpdateFunction.
|
Modifier and Type | Method and Description |
---|---|
Tuple2<K,VV> |
VertexToTuple2Map.map(Vertex<K,VV> vertex) |
Modifier and Type | Method and Description |
---|---|
Vertex<K,VV> |
Tuple2ToVertexMap.map(Tuple2<K,VV> vertex) |
Modifier and Type | Method and Description |
---|---|
TypeInformation<Tuple2<KEYOUT,VALUEOUT>> |
HadoopReduceFunction.getProducedType() |
TypeInformation<Tuple2<KEYOUT,VALUEOUT>> |
HadoopMapFunction.getProducedType() |
TypeInformation<Tuple2<KEYOUT,VALUEOUT>> |
HadoopReduceCombineFunction.getProducedType() |
Modifier and Type | Method and Description |
---|---|
void |
HadoopMapFunction.flatMap(Tuple2<KEYIN,VALUEIN> value,
Collector<Tuple2<KEYOUT,VALUEOUT>> out) |
Modifier and Type | Method and Description |
---|---|
void |
HadoopReduceCombineFunction.combine(Iterable<Tuple2<KEYIN,VALUEIN>> values,
Collector<Tuple2<KEYIN,VALUEIN>> out) |
void |
HadoopReduceCombineFunction.combine(Iterable<Tuple2<KEYIN,VALUEIN>> values,
Collector<Tuple2<KEYIN,VALUEIN>> out) |
void |
HadoopMapFunction.flatMap(Tuple2<KEYIN,VALUEIN> value,
Collector<Tuple2<KEYOUT,VALUEOUT>> out) |
void |
HadoopReduceFunction.reduce(Iterable<Tuple2<KEYIN,VALUEIN>> values,
Collector<Tuple2<KEYOUT,VALUEOUT>> out) |
void |
HadoopReduceFunction.reduce(Iterable<Tuple2<KEYIN,VALUEIN>> values,
Collector<Tuple2<KEYOUT,VALUEOUT>> out) |
void |
HadoopReduceCombineFunction.reduce(Iterable<Tuple2<KEYIN,VALUEIN>> values,
Collector<Tuple2<KEYOUT,VALUEOUT>> out) |
void |
HadoopReduceCombineFunction.reduce(Iterable<Tuple2<KEYIN,VALUEIN>> values,
Collector<Tuple2<KEYOUT,VALUEOUT>> out) |
Modifier and Type | Method and Description |
---|---|
void |
HadoopTupleUnwrappingIterator.set(Iterator<Tuple2<KEY,VALUE>> iterator)
Set the Flink iterator to wrap.
|
void |
HadoopOutputCollector.setFlinkCollector(Collector<Tuple2<KEY,VALUE>> flinkCollector)
Set the wrapped Flink collector.
|
Modifier and Type | Method and Description |
---|---|
Tuple2<byte[],byte[]> |
NestedKeyDiscarder.map(IN value) |
Modifier and Type | Method and Description |
---|---|
byte[] |
KeyDiscarder.map(Tuple2<T,byte[]> value) |
Modifier and Type | Method and Description |
---|---|
T |
RemoveRangeIndex.map(Tuple2<Integer,T> value) |
Modifier and Type | Method and Description |
---|---|
void |
AssignRangeIndex.mapPartition(Iterable<IN> values,
Collector<Tuple2<Integer,IN>> out) |
Modifier and Type | Method and Description |
---|---|
static <T> ArrayDeque<Tuple2<Long,List<T>>> |
SerializedCheckpointData.toDeque(SerializedCheckpointData[] data,
TypeSerializer<T> serializer)
De-serializes an array of SerializedCheckpointData back into an ArrayDeque of element checkpoints.
|
Modifier and Type | Method and Description |
---|---|
static <T> SerializedCheckpointData[] |
SerializedCheckpointData.fromDeque(ArrayDeque<Tuple2<Long,List<T>>> checkpoints,
TypeSerializer<T> serializer)
Converts a list of checkpoints with elements into an array of SerializedCheckpointData.
|
static <T> SerializedCheckpointData[] |
SerializedCheckpointData.fromDeque(ArrayDeque<Tuple2<Long,List<T>>> checkpoints,
TypeSerializer<T> serializer,
DataOutputSerializer outputBuffer)
Converts a list of checkpoints into an array of SerializedCheckpointData.
|
Modifier and Type | Method and Description |
---|---|
protected Tuple2<JobGraph,ClassLoader> |
JarActionHandler.getJobGraphAndClassLoader(Map<String,String> pathParams,
Map<String,String> queryParams) |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<StateHandle<T>,String>> |
ZooKeeperStateHandleStore.getAll()
Gets all available state handles from ZooKeeper.
|
List<Tuple2<StateHandle<T>,String>> |
ZooKeeperStateHandleStore.getAllSortedByName()
Gets all available state handles from ZooKeeper sorted by name (ascending).
|
Modifier and Type | Method and Description |
---|---|
void |
SpoutSourceWordCount.Tokenizer.flatMap(String value,
Collector<Tuple2<String,Integer>> out) |
Constructor and Description |
---|
DirectedOutput(List<OutputSelector<OUT>> outputSelectors,
List<Tuple2<Output<StreamRecord<OUT>>,StreamEdge>> outputs) |
Modifier and Type | Method and Description |
---|---|
<T0,T1> SingleOutputStreamOperator<Tuple2<T0,T1>> |
StreamProjection.projectTuple2()
Projects a
Tuple DataStream to the previously selected fields. |
Modifier and Type | Field and Description |
---|---|
protected Deque<Tuple2<Long,List<SessionId>>> |
MultipleIdsMessageAcknowledgingSourceBase.sessionIdsPerSnapshot |
Modifier and Type | Method and Description |
---|---|
Tuple2<StreamNode,StreamNode> |
StreamGraph.createIterationSourceAndSink(int loopId,
int sourceId,
int sinkId,
long timeout,
int parallelism) |
Modifier and Type | Method and Description |
---|---|
Set<Tuple2<StreamNode,StreamNode>> |
StreamGraph.getIterationSourceSinkPairs() |
Set<Tuple2<Integer,StreamOperator<?>>> |
StreamGraph.getOperators() |
Modifier and Type | Method and Description |
---|---|
Writer<Tuple2<K,V>> |
SequenceFileWriter.duplicate() |
Modifier and Type | Method and Description |
---|---|
void |
SequenceFileWriter.write(Tuple2<K,V> element) |
Modifier and Type | Method and Description |
---|---|
Tuple2<Tuple2<Integer,Integer>,Integer> |
IterateExample.OutputMap.map(Tuple5<Integer,Integer,Integer,Integer,Integer> value) |
Modifier and Type | Method and Description |
---|---|
Tuple2<Tuple2<Integer,Integer>,Integer> |
IterateExample.OutputMap.map(Tuple5<Integer,Integer,Integer,Integer,Integer> value) |
Modifier and Type | Method and Description |
---|---|
Tuple5<Integer,Integer,Integer,Integer,Integer> |
IterateExample.InputMap.map(Tuple2<Integer,Integer> value) |
Modifier and Type | Method and Description |
---|---|
void |
TwitterStream.SelectEnglishAndTokenizeFlatMap.flatMap(String value,
Collector<Tuple2<String,Integer>> out)
Select the language from the incoming JSON text
|
Modifier and Type | Method and Description |
---|---|
Tuple2<Long,Long> |
GroupedProcessingTimeWindowExample.SummingReducer.reduce(Tuple2<Long,Long> value1,
Tuple2<Long,Long> value2) |
Modifier and Type | Method and Description |
---|---|
Tuple2<Long,Long> |
GroupedProcessingTimeWindowExample.SummingReducer.reduce(Tuple2<Long,Long> value1,
Tuple2<Long,Long> value2) |
Tuple2<Long,Long> |
GroupedProcessingTimeWindowExample.SummingReducer.reduce(Tuple2<Long,Long> value1,
Tuple2<Long,Long> value2) |
Modifier and Type | Method and Description |
---|---|
void |
GroupedProcessingTimeWindowExample.SummingWindowFunction.apply(Long key,
Window window,
Iterable<Tuple2<Long,Long>> values,
Collector<Tuple2<Long,Long>> out) |
void |
GroupedProcessingTimeWindowExample.SummingWindowFunction.apply(Long key,
Window window,
Iterable<Tuple2<Long,Long>> values,
Collector<Tuple2<Long,Long>> out) |
Modifier and Type | Method and Description |
---|---|
void |
WordCount.Tokenizer.flatMap(String value,
Collector<Tuple2<String,Integer>> out) |
Modifier and Type | Method and Description |
---|---|
Tuple2<K,V> |
TypeInformationKeyValueSerializationSchema.deserialize(byte[] messageKey,
byte[] message,
String topic,
int partition,
long offset) |
Modifier and Type | Method and Description |
---|---|
TypeInformation<Tuple2<K,V>> |
TypeInformationKeyValueSerializationSchema.getProducedType() |
Modifier and Type | Method and Description |
---|---|
boolean |
TypeInformationKeyValueSerializationSchema.isEndOfStream(Tuple2<K,V> nextElement)
This schema never considers an element to signal end-of-stream, so this method returns always false.
|
byte[] |
TypeInformationKeyValueSerializationSchema.serializeKey(Tuple2<K,V> element) |
byte[] |
TypeInformationKeyValueSerializationSchema.serializeValue(Tuple2<K,V> element) |
Modifier and Type | Method and Description |
---|---|
static void |
ConnectedComponentsData.checkOddEvenResult(List<Tuple2<Long,Long>> lines) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.