| Package | Description |
|---|---|
| org.apache.paimon.flink.source | |
| org.apache.paimon.flink.source.align |
| Modifier and Type | Method and Description |
|---|---|
static void |
FlinkRecordsWithSplitIds.emitRecord(org.apache.flink.api.connector.source.SourceReaderContext context,
org.apache.flink.connector.file.src.reader.BulkFormat.RecordIterator<org.apache.flink.table.data.RowData> element,
org.apache.flink.api.connector.source.SourceOutput<org.apache.flink.table.data.RowData> output,
FileStoreSourceSplitState state,
FileStoreSourceReaderMetrics metrics) |
| Constructor and Description |
|---|
FileStoreSourceReader(org.apache.flink.api.connector.source.SourceReaderContext readerContext,
TableRead tableRead,
FileStoreSourceReaderMetrics metrics,
IOManager ioManager,
Long limit) |
FileStoreSourceReader(org.apache.flink.api.connector.source.SourceReaderContext readerContext,
TableRead tableRead,
FileStoreSourceReaderMetrics metrics,
IOManager ioManager,
Long limit,
org.apache.flink.connector.base.source.reader.synchronization.FutureCompletingBlockingQueue<org.apache.flink.connector.base.source.reader.RecordsWithSplitIds<org.apache.flink.connector.file.src.reader.BulkFormat.RecordIterator<org.apache.flink.table.data.RowData>>> elementsQueue) |
FileStoreSourceSplitReader(TableRead tableRead,
RecordLimiter limiter,
FileStoreSourceReaderMetrics metrics) |
| Constructor and Description |
|---|
AlignedSourceReader(org.apache.flink.api.connector.source.SourceReaderContext readerContext,
TableRead tableRead,
FileStoreSourceReaderMetrics metrics,
IOManager ioManager,
Long limit,
org.apache.flink.connector.base.source.reader.synchronization.FutureCompletingBlockingQueue<org.apache.flink.connector.base.source.reader.RecordsWithSplitIds<org.apache.flink.connector.file.src.reader.BulkFormat.RecordIterator<org.apache.flink.table.data.RowData>>> elementsQueue) |
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.