Modifier and Type | Method and Description |
---|---|
SourceReader<T,SplitT> |
Source.createReader(SourceReaderContext readerContext)
Creates a new reader to read data from the splits it gets assigned.
|
Modifier and Type | Method and Description |
---|---|
SourceReader<Long,NumberSequenceSource.NumberSequenceSplit> |
NumberSequenceSource.createReader(SourceReaderContext readerContext) |
Constructor and Description |
---|
IteratorSourceReader(SourceReaderContext context) |
Modifier and Type | Method and Description |
---|---|
SourceReader<T,HybridSourceSplit> |
HybridSource.createReader(SourceReaderContext readerContext) |
Constructor and Description |
---|
HybridSourceReader(SourceReaderContext readerContext) |
Modifier and Type | Field and Description |
---|---|
protected SourceReaderContext |
SourceReaderBase.context
The context of this source reader.
|
Modifier and Type | Method and Description |
---|---|
SourceReader<T,SplitT> |
AbstractFileSource.createReader(SourceReaderContext readerContext) |
Constructor and Description |
---|
FileSourceReader(SourceReaderContext readerContext,
BulkFormat<T,SplitT> readerFormat,
Configuration config) |
Modifier and Type | Method and Description |
---|---|
SourceReader<OUT,KafkaPartitionSplit> |
KafkaSource.createReader(SourceReaderContext readerContext) |
Constructor and Description |
---|
KafkaPartitionSplitReader(Properties props,
SourceReaderContext context,
KafkaSourceReaderMetrics kafkaSourceReaderMetrics) |
KafkaSourceReader(FutureCompletingBlockingQueue<RecordsWithSplitIds<org.apache.kafka.clients.consumer.ConsumerRecord<byte[],byte[]>>> elementsQueue,
KafkaSourceFetcherManager kafkaSourceFetcherManager,
RecordEmitter<org.apache.kafka.clients.consumer.ConsumerRecord<byte[],byte[]>,T,KafkaPartitionSplitState> recordEmitter,
Configuration config,
SourceReaderContext context,
KafkaSourceReaderMetrics kafkaSourceReaderMetrics) |
Modifier and Type | Method and Description |
---|---|
SourceReader<OUT,PulsarPartitionSplit> |
PulsarSource.createReader(SourceReaderContext readerContext) |
Modifier and Type | Method and Description |
---|---|
static <OUT> SourceReader<OUT,PulsarPartitionSplit> |
PulsarSourceReaderFactory.create(SourceReaderContext readerContext,
PulsarDeserializationSchema<OUT> deserializationSchema,
SourceConfiguration sourceConfiguration) |
Constructor and Description |
---|
PulsarDeserializationSchemaInitializationContext(SourceReaderContext readerContext) |
Constructor and Description |
---|
PulsarOrderedSourceReader(FutureCompletingBlockingQueue<RecordsWithSplitIds<PulsarMessage<OUT>>> elementsQueue,
java.util.function.Supplier<PulsarOrderedPartitionSplitReader<OUT>> splitReaderSupplier,
SourceReaderContext context,
SourceConfiguration sourceConfiguration,
org.apache.pulsar.client.api.PulsarClient pulsarClient,
org.apache.pulsar.client.admin.PulsarAdmin pulsarAdmin) |
PulsarUnorderedSourceReader(FutureCompletingBlockingQueue<RecordsWithSplitIds<PulsarMessage<OUT>>> elementsQueue,
java.util.function.Supplier<PulsarUnorderedPartitionSplitReader<OUT>> splitReaderSupplier,
SourceReaderContext context,
SourceConfiguration sourceConfiguration,
org.apache.pulsar.client.api.PulsarClient pulsarClient,
org.apache.pulsar.client.admin.PulsarAdmin pulsarAdmin,
org.apache.pulsar.client.api.transaction.TransactionCoordinatorClient coordinatorClient) |
Modifier and Type | Method and Description |
---|---|
SourceReader<OUT,FromElementsSplit> |
FromElementsSource.createReader(SourceReaderContext readerContext) |
Constructor and Description |
---|
FromElementsSourceReader(Integer limitedNum,
List<T> elements,
Boundedness boundedness,
SourceReaderContext context) |
Modifier and Type | Class and Description |
---|---|
class |
TestingReaderContext
A testing implementation of the
SourceReaderContext . |
Constructor and Description |
---|
SourceOperator(FunctionWithException<SourceReaderContext,SourceReader<OUT,SplitT>,Exception> readerFactory,
OperatorEventGateway operatorEventGateway,
SimpleVersionedSerializer<SplitT> splitSerializer,
WatermarkStrategy<OUT> watermarkStrategy,
ProcessingTimeService timeService,
Configuration configuration,
String localHostname,
boolean emitProgressiveWatermarks) |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.