Package | Description |
---|---|
org.apache.flink.connector.kafka.source.reader |
Modifier and Type | Method and Description |
---|---|
protected KafkaPartitionSplitState |
KafkaSourceReader.initializedState(KafkaPartitionSplit split) |
Modifier and Type | Method and Description |
---|---|
void |
KafkaRecordEmitter.emitRecord(Tuple3<T,Long,Long> element,
SourceOutput<T> output,
KafkaPartitionSplitState splitState) |
protected KafkaPartitionSplit |
KafkaSourceReader.toSplitType(String splitId,
KafkaPartitionSplitState splitState) |
Modifier and Type | Method and Description |
---|---|
protected void |
KafkaSourceReader.onSplitFinished(Map<String,KafkaPartitionSplitState> finishedSplitIds) |
Constructor and Description |
---|
KafkaSourceReader(FutureCompletingBlockingQueue<RecordsWithSplitIds<Tuple3<T,Long,Long>>> elementsQueue,
KafkaSourceFetcherManager<T> kafkaSourceFetcherManager,
RecordEmitter<Tuple3<T,Long,Long>,T,KafkaPartitionSplitState> recordEmitter,
Configuration config,
SourceReaderContext context,
KafkaSourceReaderMetrics kafkaSourceReaderMetrics) |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.