Modifier and Type | Method and Description |
---|---|
void |
SpoutWrapper.run(SourceFunction.SourceContext<OUT> ctx) |
Modifier and Type | Method and Description |
---|---|
void |
ContinuousFileMonitoringFunction.run(SourceFunction.SourceContext<FileInputSplit> context) |
void |
StatefulSequenceSource.run(SourceFunction.SourceContext<Long> ctx) |
void |
InputFormatSourceFunction.run(SourceFunction.SourceContext<OUT> ctx) |
void |
SocketTextStreamFunction.run(SourceFunction.SourceContext<String> ctx) |
void |
SourceFunction.run(SourceFunction.SourceContext<T> ctx)
Starts the source.
|
void |
FromSplittableIteratorFunction.run(SourceFunction.SourceContext<T> ctx) |
void |
FromIteratorFunction.run(SourceFunction.SourceContext<T> ctx) |
void |
FromElementsFunction.run(SourceFunction.SourceContext<T> ctx) |
void |
FileMonitoringFunction.run(SourceFunction.SourceContext<Tuple3<String,Long,Long>> ctx)
Deprecated.
|
Modifier and Type | Class and Description |
---|---|
static class |
StreamSource.AutomaticWatermarkContext<T>
SourceFunction.SourceContext to be used for sources with automatic timestamps
and watermark emission. |
static class |
StreamSource.ManualWatermarkContext<T>
A SourceContext for event time.
|
static class |
StreamSource.NonTimestampContext<T>
A source context that attached
-1 as a timestamp to all records, and that
does not forward watermarks. |
Modifier and Type | Method and Description |
---|---|
<T> DataStream<T> |
StreamExecutionEnvironment.addSource(scala.Function1<SourceFunction.SourceContext<T>,scala.runtime.BoxedUnit> function,
TypeInformation<T> evidence$9)
Create a DataStream using a user defined source function for arbitrary
source functionality.
|
Modifier and Type | Method and Description |
---|---|
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer09.createFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> thisSubtaskPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext) |
protected abstract AbstractFetcher<T,?> |
FlinkKafkaConsumerBase.createFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> thisSubtaskPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext)
Creates the fetcher that connect to the Kafka brokers, pulls data, deserialized the
data, and emits it into the data streams.
|
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer08.createFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> thisSubtaskPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext) |
void |
FlinkKafkaConsumerBase.run(SourceFunction.SourceContext<T> sourceContext) |
Constructor and Description |
---|
Kafka09Fetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> assignedPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext,
KeyedDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
boolean useMetrics) |
Constructor and Description |
---|
AbstractFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> assignedPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext,
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) |
Modifier and Type | Method and Description |
---|---|
void |
NiFiSource.run(SourceFunction.SourceContext<NiFiDataPacket> ctx) |
Modifier and Type | Method and Description |
---|---|
void |
RMQSource.run(SourceFunction.SourceContext<OUT> ctx) |
Modifier and Type | Method and Description |
---|---|
void |
TwitterSource.run(SourceFunction.SourceContext<String> ctx) |
Modifier and Type | Method and Description |
---|---|
void |
WikipediaEditsSource.run(SourceFunction.SourceContext ctx) |
Modifier and Type | Method and Description |
---|---|
void |
IncrementalLearningSkeleton.FiniteNewDataSource.run(SourceFunction.SourceContext<Integer> ctx) |
void |
IncrementalLearningSkeleton.FiniteTrainingDataSource.run(SourceFunction.SourceContext<Integer> collector) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.