Modifier and Type | Class and Description |
---|---|
class |
HBaseUpsertSinkFunction
The upsert sink for HBase.
|
Modifier and Type | Class and Description |
---|---|
class |
StateBootstrapFunction<IN>
Interface for writing elements to operator state.
|
Modifier and Type | Class and Description |
---|---|
class |
TwoPhaseCommitSinkFunction<IN,TXN,CONTEXT>
This is a recommended base class for all of the
SinkFunction that intend to implement exactly-once semantic. |
Modifier and Type | Class and Description |
---|---|
class |
StreamingFileSink<IN>
Sink that emits its input elements to
FileSystem files within buckets. |
Modifier and Type | Class and Description |
---|---|
class |
ContinuousFileMonitoringFunction<OUT>
This is the single (non-parallel) monitoring task which takes a
FileInputFormat
and, depending on the FileProcessingMode and the FilePathFilter , it is responsible for:
Monitoring a user-provided path.
Deciding which files should be further read and processed.
Creating the splits corresponding to those files.
Assigning them to downstream tasks for further processing.
|
class |
FromElementsFunction<T>
A stream source function that returns a sequence of elements.
|
class |
MessageAcknowledgingSourceBase<Type,UId>
Abstract base class for data sources that receive elements from a message queue and
acknowledge them back by IDs.
|
class |
MultipleIdsMessageAcknowledgingSourceBase<Type,UId,SessionId>
Abstract base class for data sources that receive elements from a message queue and
acknowledge them back by IDs.
|
class |
StatefulSequenceSource
A stateful streaming source that emits each number from a given interval exactly once,
possibly in parallel.
|
Modifier and Type | Class and Description |
---|---|
class |
AbstractCassandraTupleSink<IN>
Abstract sink to write tuple-like values into a Cassandra cluster.
|
class |
CassandraPojoSink<IN>
Flink Sink to save data into a Cassandra cluster using
Mapper,
which it uses annotations from
com.datastax.driver.mapping.annotations.
|
class |
CassandraRowSink
A SinkFunction to write Row records into a Cassandra table.
|
class |
CassandraScalaProductSink<IN extends scala.Product>
Sink to write scala tuples and case classes into a Cassandra cluster.
|
class |
CassandraSinkBase<IN,V>
CassandraSinkBase is the common abstract class of
CassandraPojoSink and CassandraTupleSink . |
class |
CassandraTupleSink<IN extends Tuple>
Sink to write Flink
Tuple s into a Cassandra cluster. |
Modifier and Type | Class and Description |
---|---|
class |
ElasticsearchSinkBase<T,C extends AutoCloseable>
Base class for all Flink Elasticsearch Sinks.
|
Modifier and Type | Class and Description |
---|---|
class |
ElasticsearchSink<T>
Elasticsearch 2.x sink that requests multiple
ActionRequests
against a cluster for each incoming element. |
Modifier and Type | Class and Description |
---|---|
class |
BucketingSink<T>
Deprecated.
Please use the
StreamingFileSink
instead. |
Modifier and Type | Class and Description |
---|---|
class |
PubSubSink<IN>
A sink function that outputs to PubSub.
|
Modifier and Type | Class and Description |
---|---|
class |
FlinkKafkaConsumer<T>
The Flink Kafka Consumer is a streaming data source that pulls a parallel data stream from
Apache Kafka.
|
class |
FlinkKafkaConsumer010<T>
The Flink Kafka Consumer is a streaming data source that pulls a parallel data stream from
Apache Kafka 0.10.x.
|
class |
FlinkKafkaConsumer011<T>
The Flink Kafka Consumer is a streaming data source that pulls a parallel data stream from
Apache Kafka 0.11.x.
|
class |
FlinkKafkaConsumer08<T>
The Flink Kafka Consumer is a streaming data source that pulls a parallel data stream from
Apache Kafka 0.8.x.
|
class |
FlinkKafkaConsumer081<T>
Deprecated.
|
class |
FlinkKafkaConsumer082<T>
Deprecated.
|
class |
FlinkKafkaConsumer09<T>
The Flink Kafka Consumer is a streaming data source that pulls a parallel data stream from
Apache Kafka 0.9.x.
|
class |
FlinkKafkaConsumerBase<T>
Base class of all Flink Kafka Consumer data sources.
|
class |
FlinkKafkaProducer<IN>
Flink Sink to produce data into a Kafka topic.
|
class |
FlinkKafkaProducer010<T>
Flink Sink to produce data into a Kafka topic.
|
class |
FlinkKafkaProducer011<IN>
Flink Sink to produce data into a Kafka topic.
|
class |
FlinkKafkaProducer08<IN>
Flink Sink to produce data into a Kafka topic.
|
class |
FlinkKafkaProducer09<IN>
Flink Sink to produce data into a Kafka topic.
|
class |
FlinkKafkaProducerBase<IN>
Flink Sink to produce data into a Kafka topic.
|
Modifier and Type | Class and Description |
---|---|
class |
FlinkDynamoDBStreamsConsumer<T>
Consume events from DynamoDB streams.
|
class |
FlinkKinesisConsumer<T>
The Flink Kinesis Consumer is an exactly-once parallel streaming data source that subscribes to multiple AWS Kinesis
streams within the same AWS service region, and can handle resharding of streams.
|
class |
FlinkKinesisProducer<OUT>
The FlinkKinesisProducer allows to produce from a Flink DataStream into Kinesis.
|
Modifier and Type | Class and Description |
---|---|
class |
RMQSource<OUT>
RabbitMQ source (consumer) which reads from a queue and acknowledges messages on checkpoints.
|
Modifier and Type | Class and Description |
---|---|
class |
SequenceGeneratorSource
This source function generates a sequence of long values per key.
|
Modifier and Type | Class and Description |
---|---|
class |
ArtificalOperatorStateMapper<IN,OUT>
A self-verifiable
RichMapFunction used to verify checkpointing and restore semantics for various
kinds of operator state. |
class |
ArtificialKeyedStateMapper<IN,OUT>
A generic, stateful
MapFunction that allows specifying what states to maintain
based on a provided list of ArtificialStateBuilder s. |
Modifier and Type | Class and Description |
---|---|
class |
UpdatableTopNFunction
The function could handle update input stream.
|
Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.