Package | Description |
---|---|
org.apache.flink.streaming.api.datastream | |
org.apache.flink.streaming.api.functions | |
org.apache.flink.streaming.api.functions.timestamps | |
org.apache.flink.streaming.api.scala | |
org.apache.flink.streaming.connectors.kafka | |
org.apache.flink.streaming.connectors.kafka.internal | |
org.apache.flink.streaming.connectors.kafka.internals | |
org.apache.flink.streaming.runtime.operators |
This package contains the operators that perform the stream transformations.
|
Modifier and Type | Method and Description |
---|---|
SingleOutputStreamOperator<T> |
DataStream.assignTimestampsAndWatermarks(AssignerWithPeriodicWatermarks<T> timestampAndWatermarkAssigner)
Assigns timestamps to the elements in the data stream and periodically creates
watermarks to signal event time progress.
|
Modifier and Type | Class and Description |
---|---|
class |
IngestionTimeExtractor<T>
A timestamp assigner that assigns timestamps based on the machine's wall clock.
|
Modifier and Type | Class and Description |
---|---|
class |
AscendingTimestampExtractor<T>
A timestamp assigner and watermark generator for streams where timestamps are monotonously
ascending.
|
class |
BoundedOutOfOrdernessTimestampExtractor<T>
This is a
AssignerWithPeriodicWatermarks used to emit Watermarks that lag behind the element with
the maximum timestamp (in event time) seen so far by a fixed amount of time, t_late . |
Modifier and Type | Method and Description |
---|---|
DataStream<T> |
DataStream.assignTimestampsAndWatermarks(AssignerWithPeriodicWatermarks<T> assigner)
Assigns timestamps to the elements in the data stream and periodically creates
watermarks to signal event time progress.
|
Modifier and Type | Method and Description |
---|---|
FlinkKafkaConsumerBase<T> |
FlinkKafkaConsumerBase.assignTimestampsAndWatermarks(AssignerWithPeriodicWatermarks<T> assigner)
Specifies an
AssignerWithPunctuatedWatermarks to emit watermarks in a punctuated manner. |
Constructor and Description |
---|
Kafka010Fetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> assignedPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
boolean enableCheckpointing,
String taskNameWithSubtasks,
MetricGroup metricGroup,
KeyedDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
boolean useMetrics) |
Kafka09Fetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> assignedPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
boolean enableCheckpointing,
String taskNameWithSubtasks,
MetricGroup metricGroup,
KeyedDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
boolean useMetrics) |
Constructor and Description |
---|
KafkaTopicPartitionStateWithPeriodicWatermarks(KafkaTopicPartition partition,
KPH kafkaPartitionHandle,
AssignerWithPeriodicWatermarks<T> timestampsAndWatermarks) |
Constructor and Description |
---|
AbstractFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> assignedPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
boolean useMetrics) |
Kafka08Fetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> assignedPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext,
KeyedDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long invalidOffsetBehavior,
long autoCommitInterval,
boolean useMetrics) |
Constructor and Description |
---|
TimestampsAndPeriodicWatermarksOperator(AssignerWithPeriodicWatermarks<T> assigner) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.