Modifier and Type | Method and Description |
---|---|
Tuple2<Long,Long> |
BinaryInputFormat.getCurrentState() |
Modifier and Type | Method and Description |
---|---|
void |
BinaryInputFormat.reopen(FileInputSplit split,
Tuple2<Long,Long> state) |
Modifier and Type | Method and Description |
---|---|
Tuple2<TypeSerializer<?>,TypeSerializerConfigSnapshot> |
CompositeTypeSerializerConfigSnapshot.getSingleNestedSerializerAndConfig() |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<TypeSerializer<?>,TypeSerializerConfigSnapshot>> |
CompositeTypeSerializerConfigSnapshot.getNestedSerializersAndConfigs() |
static List<Tuple2<TypeSerializer<?>,TypeSerializerConfigSnapshot>> |
TypeSerializerSerializationUtil.readSerializersAndConfigsWithResilience(DataInputView in,
ClassLoader userCodeClassLoader)
Reads from a data input view a list of serializers and their corresponding config snapshots
written using
TypeSerializerSerializationUtil.writeSerializersAndConfigsWithResilience(DataOutputView, List) . |
Modifier and Type | Method and Description |
---|---|
static void |
TypeSerializerSerializationUtil.writeSerializersAndConfigsWithResilience(DataOutputView out,
List<Tuple2<TypeSerializer<?>,TypeSerializerConfigSnapshot>> serializersAndConfigs)
Write a list of serializers and their corresponding config snapshots to the provided
data output view.
|
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)
Deprecated.
Please use
org.apache.flink.hadoopcompatibility.HadoopInputs#createHadoopInput(org.apache.hadoop.mapreduce.InputFormat
from the flink-hadoop-compatibility module. |
<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)
Deprecated.
Please use
org.apache.flink.hadoopcompatibility.HadoopInputs#createHadoopInput(org.apache.hadoop.mapred.InputFormat
from the flink-hadoop-compatibility module. |
<K,V> DataSource<Tuple2<K,V>> |
ExecutionEnvironment.readHadoopFile(org.apache.hadoop.mapred.FileInputFormat<K,V> mapredInputFormat,
Class<K> key,
Class<V> value,
String inputPath)
Deprecated.
Please use
org.apache.flink.hadoopcompatibility.HadoopInputs#readHadoopFile(org.apache.hadoop.mapred.FileInputFormat
from the flink-hadoop-compatibility module. |
<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)
Deprecated.
Please use
org.apache.flink.hadoopcompatibility.HadoopInputs#readHadoopFile(org.apache.hadoop.mapreduce.lib.input.FileInputFormat
from the flink-hadoop-compatibility module. |
<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)
Deprecated.
Please use
org.apache.flink.hadoopcompatibility.HadoopInputs#readHadoopFile(org.apache.hadoop.mapreduce.lib.input.FileInputFormat
from the flink-hadoop-compatibility module. |
<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)
Deprecated.
Please use
org.apache.flink.hadoopcompatibility.HadoopInputs#readHadoopFile(org.apache.hadoop.mapred.FileInputFormat
from the flink-hadoop-compatibility module. |
<K,V> DataSource<Tuple2<K,V>> |
ExecutionEnvironment.readSequenceFile(Class<K> key,
Class<V> value,
String inputPath)
Deprecated.
Please use
org.apache.flink.hadoopcompatibility.HadoopInputs#readSequenceFile(Class
from the flink-hadoop-compatibility module. |
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 |
---|---|
Tuple2<Long,Long> |
AvroInputFormat.getCurrentState() |
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 |
---|---|
void |
AvroInputFormat.reopen(FileInputSplit split,
Tuple2<Long,Long> state) |
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 |
---|---|
LinkedHashMap<String,Tuple2<TypeSerializer<?>,TypeSerializerConfigSnapshot>> |
PojoSerializer.PojoSerializerConfigSnapshot.getFieldToSerializerConfigSnapshot() |
HashMap<Class<?>,Tuple2<TypeSerializer<?>,TypeSerializerConfigSnapshot>> |
PojoSerializer.PojoSerializerConfigSnapshot.getNonRegisteredSubclassesToSerializerConfigSnapshots() |
LinkedHashMap<Class<?>,Tuple2<TypeSerializer<?>,TypeSerializerConfigSnapshot>> |
PojoSerializer.PojoSerializerConfigSnapshot.getRegisteredSubclassesToSerializerConfigSnapshots() |
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 |
---|---|
Tuple2<Collection<Map<String,List<T>>>,Collection<Tuple2<Map<String,List<T>>,Long>>> |
NFA.process(T event,
long timestamp)
Processes the next input event.
|
Modifier and Type | Method and Description |
---|---|
Tuple2<Collection<Map<String,List<T>>>,Collection<Tuple2<Map<String,List<T>>,Long>>> |
NFA.process(T event,
long timestamp)
Processes the next input event.
|
Modifier and Type | Method and Description |
---|---|
static <K,T> SingleOutputStreamOperator<Either<Tuple2<Map<String,List<T>>,Long>,Map<String,List<T>>>> |
CEPOperatorUtils.createTimeoutPatternStream(DataStream<T> inputStream,
Pattern<T,?> pattern)
Creates a data stream containing fully matching event patterns or partially matching event
patterns which have timed out.
|
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,LongValue>> |
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,LongValue>> |
Graph.inDegrees()
Return the in-degree of all vertices in the graph
|
DataSet<Tuple2<K,LongValue>> |
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 |
---|---|
Edge<K,Tuple2<EV,D>> |
DegreeAnnotationFunctions.JoinEdgeWithVertexDegree.join(Edge<K,EV> edge,
Vertex<K,D> vertex) |
Modifier and Type | Method and Description |
---|---|
Edge<K,Tuple3<EV,D,D>> |
DegreeAnnotationFunctions.JoinEdgeDegreeWithVertexDegree.join(Edge<K,Tuple2<EV,D>> edge,
Vertex<K,D> vertex) |
Modifier and Type | Method and Description |
---|---|
DataSet<Edge<K,Tuple2<EV,VertexDegrees.Degrees>>> |
EdgeTargetDegrees.runInternal(Graph<K,VV,EV> input) |
DataSet<Edge<K,Tuple2<EV,VertexDegrees.Degrees>>> |
EdgeSourceDegrees.runInternal(Graph<K,VV,EV> input) |
Modifier and Type | Method and Description |
---|---|
DataSet<Edge<K,Tuple2<EV,LongValue>>> |
EdgeTargetDegree.runInternal(Graph<K,VV,EV> input) |
DataSet<Edge<K,Tuple2<EV,LongValue>>> |
EdgeSourceDegree.runInternal(Graph<K,VV,EV> input) |
Modifier and Type | Method and Description |
---|---|
Graph<KB,VVB,Tuple2<EV,EV>> |
BipartiteGraph.projectionBottomSimple()
Convert a bipartite graph into an undirected graph that contains only bottom vertices.
|
Graph<KT,VVT,Tuple2<EV,EV>> |
BipartiteGraph.projectionTopSimple()
Convert a bipartite graph into an undirected graph that contains only top vertices.
|
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 | Class and Description |
---|---|
static class |
PageRank.Result<T>
Wraps the
Tuple2 to encapsulate results from the PageRank algorithm. |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<String,DataSet<?>>> |
VertexCentricConfiguration.getBcastVars()
Get the broadcast variables of the compute function.
|
TypeInformation<Tuple2<K,Either<NullValue,Message>>> |
VertexCentricIteration.MessageCombinerUdf.getProducedType() |
Modifier and Type | Method and Description |
---|---|
void |
VertexCentricIteration.MessageCombinerUdf.combine(Iterable<Tuple2<K,Either<NullValue,Message>>> values,
Collector<Tuple2<K,Either<NullValue,Message>>> out) |
void |
VertexCentricIteration.MessageCombinerUdf.combine(Iterable<Tuple2<K,Either<NullValue,Message>>> values,
Collector<Tuple2<K,Either<NullValue,Message>>> out) |
void |
VertexCentricIteration.MessageCombinerUdf.reduce(Iterable<Tuple2<K,Either<NullValue,Message>>> messages,
Collector<Tuple2<K,Either<NullValue,Message>>> out) |
void |
VertexCentricIteration.MessageCombinerUdf.reduce(Iterable<Tuple2<K,Either<NullValue,Message>>> messages,
Collector<Tuple2<K,Either<NullValue,Message>>> out) |
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.getGatherBcastVars()
Get the broadcast variables of the GatherFunction.
|
List<Tuple2<String,DataSet<?>>> |
ScatterGatherConfiguration.getScatterBcastVars()
Get the broadcast variables of the ScatterFunction.
|
Modifier and Type | Method and Description |
---|---|
Tuple2<K,K> |
EdgeToTuple2Map.map(Edge<K,EV> edge) |
Tuple2<K,VV> |
VertexToTuple2Map.map(Vertex<K,VV> vertex) |
Modifier and Type | Method and Description |
---|---|
Edge<K,NullValue> |
Tuple2ToEdgeMap.map(Tuple2<K,K> tuple) |
Vertex<K,VV> |
Tuple2ToVertexMap.map(Tuple2<K,VV> tuple) |
Modifier and Type | Method and Description |
---|---|
TypeInformation<Tuple2<KEYOUT,VALUEOUT>> |
HadoopReduceFunction.getProducedType() |
TypeInformation<Tuple2<KEYOUT,VALUEOUT>> |
HadoopReduceCombineFunction.getProducedType() |
TypeInformation<Tuple2<KEYOUT,VALUEOUT>> |
HadoopMapFunction.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 |
---|---|
List<Tuple2<com.netflix.fenzo.TaskRequest,String>> |
LaunchCoordinator.Assign.tasks() |
Constructor and Description |
---|
Assign(List<Tuple2<com.netflix.fenzo.TaskRequest,String>> tasks) |
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 |
---|---|
static Tuple2<Savepoint,StreamStateHandle> |
SavepointStore.loadSavepointWithHandle(String savepointFileOrDirectory,
ClassLoader classLoader)
Loads the savepoint at the specified path.
|
Modifier and Type | Method and Description |
---|---|
static Tuple2<String,Integer> |
HighAvailabilityServicesUtils.getJobManagerAddress(Configuration configuration)
Returns the JobManager's hostname and port extracted from the given
Configuration . |
Modifier and Type | Method and Description |
---|---|
Future<Iterable<SlotOffer>> |
SlotPoolGateway.offerSlots(Iterable<Tuple2<AllocatedSlot,SlotOffer>> offers) |
Iterable<SlotOffer> |
SlotPool.offerSlots(Iterable<Tuple2<AllocatedSlot,SlotOffer>> offers) |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<String,Configuration>> |
MetricRegistryConfiguration.getReporterConfigurations() |
Constructor and Description |
---|
MetricRegistryConfiguration(ScopeFormats scopeFormats,
char delimiter,
List<Tuple2<String,Configuration>> reporterConfigurations) |
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 <K,N> Tuple2<K,N> |
KvStateRequestSerializer.deserializeKeyAndNamespace(byte[] serializedKeyAndNamespace,
TypeSerializer<K> keySerializer,
TypeSerializer<N> namespaceSerializer)
Deserializes the key and namespace into a
Tuple2 . |
Modifier and Type | Method and Description |
---|---|
Future<Tuple2<Gateway,Success>> |
RetryingRegistration.getFuture() |
Modifier and Type | Method and Description |
---|---|
Iterator<Tuple2<Integer,Long>> |
KeyGroupRangeOffsets.iterator() |
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(org.apache.flink.runtime.webmonitor.handlers.JarActionHandler.JarActionHandlerConfig config) |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<RetrievableStateHandle<T>,String>> |
ZooKeeperStateHandleStore.getAllAndLock()
Gets all available state handles from ZooKeeper and locks the respective state nodes.
|
List<Tuple2<RetrievableStateHandle<T>,String>> |
ZooKeeperStateHandleStore.getAllSortedByNameAndLock()
Gets all available state handles from ZooKeeper sorted by name (ascending) and locks the
respective state nodes.
|
Modifier and Type | Method and Description |
---|---|
Tuple2<String,Integer> |
SpoutSplitExample.Enrich.map(Integer value) |
Modifier and Type | Method and Description |
---|---|
void |
SpoutSourceWordCount.Tokenizer.flatMap(String value,
Collector<Tuple2<String,Integer>> out) |
Constructor and Description |
---|
CopyingDirectedOutput(List<OutputSelector<OUT>> outputSelectors,
List<Tuple2<Output<StreamRecord<OUT>>,StreamEdge>> outputs) |
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 List<Tuple2<String,DistributedCache.DistributedCacheEntry>> |
StreamExecutionEnvironment.cacheFile |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<String,DistributedCache.DistributedCacheEntry>> |
StreamExecutionEnvironment.getCachedFiles()
Get the list of cached files that were registered for distribution among the task managers.
|
Modifier and Type | Field and Description |
---|---|
protected ArrayDeque<Tuple2<Long,List<UId>>> |
MessageAcknowledgingSourceBase.pendingCheckpoints
The list with IDs from checkpoints that were triggered, but not yet completed or notified of
completion.
|
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,
int maxParallelism,
ResourceSpec minResources,
ResourceSpec preferredResources) |
Modifier and Type | Method and Description |
---|---|
Set<Tuple2<StreamNode,StreamNode>> |
StreamGraph.getIterationSourceSinkPairs() |
Set<Tuple2<Integer,StreamOperator<?>>> |
StreamGraph.getOperators() |
Modifier and Type | Method and Description |
---|---|
List<Tuple2<String,DistributedCache.DistributedCacheEntry>> |
StreamExecutionEnvironment.getCachedFiles()
Gets cache files.
|
Modifier and Type | Method and Description |
---|---|
Writer<Tuple2<K,V>> |
AvroKeyValueSinkWriter.duplicate() |
Writer<Tuple2<K,V>> |
SequenceFileWriter.duplicate() |
Modifier and Type | Method and Description |
---|---|
void |
AvroKeyValueSinkWriter.write(Tuple2<K,V> element) |
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 |
---|---|
Tuple2<String,Integer> |
WindowJoinSampleData.GradeSource.next() |
Tuple2<String,Integer> |
WindowJoinSampleData.SalarySource.next() |
Modifier and Type | Method and Description |
---|---|
static DataStream<Tuple2<String,Integer>> |
WindowJoinSampleData.GradeSource.getSource(StreamExecutionEnvironment env,
long rate) |
static DataStream<Tuple2<String,Integer>> |
WindowJoinSampleData.SalarySource.getSource(StreamExecutionEnvironment env,
long rate) |
Modifier and Type | Method and Description |
---|---|
static DataStream<Tuple3<String,Integer,Integer>> |
WindowJoin.runWindowJoin(DataStream<Tuple2<String,Integer>> grades,
DataStream<Tuple2<String,Integer>> salaries,
long windowSize) |
static DataStream<Tuple3<String,Integer,Integer>> |
WindowJoin.runWindowJoin(DataStream<Tuple2<String,Integer>> grades,
DataStream<Tuple2<String,Integer>> salaries,
long windowSize) |
Modifier and Type | Method and Description |
---|---|
void |
SideOutputExample.Tokenizer.processElement(String value,
ProcessFunction.Context ctx,
Collector<Tuple2<String,Integer>> out) |
Modifier and Type | Method and Description |
---|---|
void |
TwitterExample.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) |
Constructor and Description |
---|
MergingWindowSet(MergingWindowAssigner<?,W> windowAssigner,
ListState<Tuple2<W,W>> state)
Restores a
MergingWindowSet from the given state. |
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 |
---|---|
String |
TypeInformationKeyValueSerializationSchema.getTargetTopic(Tuple2<K,V> element) |
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 |
---|---|
<T> DataStream<Tuple2<Boolean,T>> |
StreamTableEnvironment.toRetractStream(Table table,
Class<T> clazz)
Converts the given
Table into a DataStream of add and retract messages. |
<T> DataStream<Tuple2<Boolean,T>> |
StreamTableEnvironment.toRetractStream(Table table,
Class<T> clazz,
StreamQueryConfig queryConfig)
Converts the given
Table into a DataStream of add and retract messages. |
<T> DataStream<Tuple2<Boolean,T>> |
StreamTableEnvironment.toRetractStream(Table table,
TypeInformation<T> typeInfo)
Converts the given
Table into a DataStream of add and retract messages. |
<T> DataStream<Tuple2<Boolean,T>> |
StreamTableEnvironment.toRetractStream(Table table,
TypeInformation<T> typeInfo,
StreamQueryConfig queryConfig)
Converts the given
Table into a DataStream of add and retract messages. |
Modifier and Type | Class and Description |
---|---|
class |
BigIntegralAvgAccumulator
The initial accumulator for Big Integral Avg aggregate function
|
class |
DecimalAvgAccumulator
The initial accumulator for Big Decimal Avg aggregate function
|
class |
DecimalSumAccumulator
The initial accumulator for Big Decimal Sum aggregate function
|
class |
DecimalSumWithRetractAccumulator
The initial accumulator for Big Decimal Sum with retract aggregate function
|
class |
FloatingAvgAccumulator
The initial accumulator for Floating Avg aggregate function
|
class |
IntegralAvgAccumulator
The initial accumulator for Integral Avg aggregate function
|
class |
MaxAccumulator<T>
The initial accumulator for Max aggregate function
|
class |
MaxWithRetractAccumulator<T>
The initial accumulator for Max with retraction aggregate function
|
class |
MinAccumulator<T>
The initial accumulator for Min aggregate function
|
class |
MinWithRetractAccumulator<T>
The initial accumulator for Min with retraction aggregate function
|
class |
SumAccumulator<T>
The initial accumulator for Sum aggregate function
|
class |
SumWithRetractAccumulator<T>
The initial accumulator for Sum with retract aggregate function
|
Modifier and Type | Method and Description |
---|---|
Tuple2<Boolean,Object> |
CRowInputJavaTupleOutputMapRunner.map(CRow in) |
Modifier and Type | Method and Description |
---|---|
TypeInformation<Tuple2<Boolean,Object>> |
CRowInputJavaTupleOutputMapRunner.getProducedType() |
TypeInformation<Tuple2<Boolean,Object>> |
CRowInputJavaTupleOutputMapRunner.returnType() |
Constructor and Description |
---|
CRowInputJavaTupleOutputMapRunner(String name,
String code,
TypeInformation<Tuple2<Boolean,Object>> returnType) |
Modifier and Type | Method and Description |
---|---|
TupleTypeInfo<Tuple2<Boolean,T>> |
UpsertStreamTableSink.getOutputType() |
TupleTypeInfo<Tuple2<Boolean,T>> |
RetractStreamTableSink.getOutputType() |
Modifier and Type | Method and Description |
---|---|
void |
UpsertStreamTableSink.emitDataStream(DataStream<Tuple2<Boolean,T>> dataStream)
Emits the DataStream.
|
void |
RetractStreamTableSink.emitDataStream(DataStream<Tuple2<Boolean,T>> dataStream)
Emits the DataStream.
|
Modifier and Type | Method and Description |
---|---|
static void |
ConnectedComponentsData.checkOddEvenResult(List<Tuple2<Long,Long>> lines) |
Copyright © 2014–2018 The Apache Software Foundation. All rights reserved.