Modifier and Type | Class and Description |
---|---|
static class |
FileSinkProgram.Generator
Data-generating source function.
|
Modifier and Type | Class and Description |
---|---|
static class |
StreamSQLTestProgram.Generator
Data-generating source function.
|
Modifier and Type | Interface and Description |
---|---|
interface |
ExternallyInducedSource<T,CD>
Sources that implement this interface delay checkpoints when receiving a trigger message from the
checkpoint coordinator to the point when their input data/events indicate that a checkpoint
should be triggered.
|
Modifier and Type | Method and Description |
---|---|
<OUT> DataStreamSource<OUT> |
StreamExecutionEnvironment.addSource(SourceFunction<OUT> function)
Adds a Data Source to the streaming topology.
|
<OUT> DataStreamSource<OUT> |
StreamExecutionEnvironment.addSource(SourceFunction<OUT> function,
String sourceName)
Adds a data source with a custom type information thus opening a
DataStream . |
<OUT> DataStreamSource<OUT> |
StreamExecutionEnvironment.addSource(SourceFunction<OUT> function,
String sourceName,
TypeInformation<OUT> typeInfo)
Ads a data source with a custom type information thus opening a
DataStream . |
<OUT> DataStreamSource<OUT> |
StreamExecutionEnvironment.addSource(SourceFunction<OUT> function,
TypeInformation<OUT> typeInfo)
Ads a data source with a custom type information thus opening a
DataStream . |
Modifier and Type | Interface and Description |
---|---|
interface |
ParallelSourceFunction<OUT>
A stream data source that is executed in parallel.
|
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. |
class |
FileMonitoringFunction
Deprecated.
Internal class deprecated in favour of
ContinuousFileMonitoringFunction . |
class |
FromElementsFunction<T>
A stream source function that returns a sequence of elements.
|
class |
FromIteratorFunction<T>
A
SourceFunction that reads elements from an Iterator and emits them. |
class |
FromSplittableIteratorFunction<T>
A
SourceFunction that reads elements from an SplittableIterator and emits them. |
class |
InputFormatSourceFunction<OUT>
A
SourceFunction that reads data using an InputFormat . |
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 |
RichParallelSourceFunction<OUT>
Base class for implementing a parallel data source.
|
class |
RichSourceFunction<OUT>
Base class for implementing a parallel data source that has access to context information (via
AbstractRichFunction.getRuntimeContext() ) and additional life-cycle methods (AbstractRichFunction.open(org.apache.flink.configuration.Configuration) and AbstractRichFunction.close() . |
class |
SocketTextStreamFunction
A source function that reads strings from a socket.
|
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 |
DataGeneratorSource<T>
A data generator source that abstract data generator.
|
Modifier and Type | Class and Description |
---|---|
class |
StreamSource<OUT,SRC extends SourceFunction<OUT>>
StreamOperator for streaming sources. |
Modifier and Type | Class and Description |
---|---|
class |
PubSubSource<OUT>
PubSub Source, this Source will consume PubSub messages from a subscription and Acknowledge them
on the next checkpoint.
|
Modifier and Type | Class and Description |
---|---|
class |
FlinkKafkaConsumer<T>
Deprecated.
|
class |
FlinkKafkaConsumerBase<T>
Base class of all Flink Kafka Consumer data sources.
|
Modifier and Type | Class and Description |
---|---|
class |
FlinkKafkaShuffleConsumer<T>
Flink Kafka Shuffle Consumer Function.
|
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.
|
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 |
WikipediaEditsSource
This class is a SourceFunction that reads
WikipediaEditEvent instances from the IRC
channel #en.wikipedia . |
Modifier and Type | Class and Description |
---|---|
class |
SimpleSource
A checkpointed source.
|
Modifier and Type | Class and Description |
---|---|
class |
EventsGeneratorSource
A event stream source that generates the events on the fly.
|
Modifier and Type | Class and Description |
---|---|
class |
CarSource
A simple in-memory source.
|
Modifier and Type | Class and Description |
---|---|
class |
SourceStreamTask<OUT,SRC extends SourceFunction<OUT>,OP extends StreamSource<OUT,SRC>>
StreamTask for executing a StreamSource . |
Modifier and Type | Class and Description |
---|---|
static class |
PeriodicStreamingJob.PeriodicSourceGenerator
Data-generating source function.
|
class |
SequenceGeneratorSource
This source function generates a sequence of long values per key.
|
Modifier and Type | Class and Description |
---|---|
class |
FiniteTestSource<T>
A stream source that: 1) emits a list of elements without allowing checkpoints, 2) then waits for
two more checkpoints to complete, 3) then re-emits the same elements before 4) waiting for
another two checkpoints and 5) exiting.
|
Modifier and Type | Method and Description |
---|---|
SourceFunction<RowData> |
SourceFunctionProvider.createSourceFunction()
Creates a
SourceFunction instance. |
Modifier and Type | Method and Description |
---|---|
static SourceFunctionProvider |
SourceFunctionProvider.of(SourceFunction<RowData> sourceFunction,
boolean isBounded)
Helper method for creating a static provider.
|
Modifier and Type | Class and Description |
---|---|
class |
SocketSourceFunction
The
SocketSourceFunction opens a socket and consumes bytes. |
Modifier and Type | Method and Description |
---|---|
protected Transformation<RowData> |
CommonExecTableSourceScan.createSourceFunctionTransformation(StreamExecutionEnvironment env,
SourceFunction<RowData> function,
boolean isBounded,
String operatorName,
TypeInformation<RowData> outputTypeInfo)
Adopted from
StreamExecutionEnvironment.addSource(SourceFunction, String,
TypeInformation) but with custom Boundedness . |
Modifier and Type | Class and Description |
---|---|
class |
ArrowSourceFunction
An Arrow
SourceFunction which takes the serialized arrow record batch data as input. |
Modifier and Type | Class and Description |
---|---|
static class |
StatefulStreamingJob.MySource
Stub source that emits one record per second.
|
Modifier and Type | Class and Description |
---|---|
class |
TransactionSource
A stream of transactions.
|
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.