Read#
Interface#
-
class TableRead#
Given a
Splitor a list ofSplit, generate a reader for batch reading.A real-time split may create at most one successful reader. Calls using the same real-time split must not run concurrently, and the split must not be reused by multiple queries. A failed reader creation may be retried sequentially while its read-view ticket remains valid.
Public Functions
-
virtual ~TableRead() = default#
Creates a
BatchReaderinstance for reading data.This method creates a BatchReader that will be responsible for reading data from the provided splits.
Note
BatchReaders created by the sameTableReadare not thread-safe for concurrent reading.- Parameters:
splits – A vector of shared pointers to
Splitinstances representing the data to be read.- Returns:
A Result containing a unique pointer to the
BatchReaderinstance.
Creates a
BatchReaderinstance for a single split.- Parameters:
split – A shared pointer to the
Splitinstance that defines the data to be read.- Returns:
A Result containing a unique pointer to the
BatchReaderinstance.
Creates a
CountReaderfor count queries on the specified splits.Implementations may override this to provide a more efficient count path.
Public Static Functions
-
static Result<std::unique_ptr<TableRead>> Create(std::unique_ptr<ReadContext> context)#
Create an instance of
TableRead.- Parameters:
context – A unique pointer to the
ReadContextused for read operations.- Returns:
A Result containing a unique pointer to the
TableReadinstance.
-
virtual ~TableRead() = default#
-
class ReadContextBuilder#
ReadContextBuilderused to build aReadContext, has input validation.Public Functions
-
explicit ReadContextBuilder(const std::string &path)#
Constructs a
ReadContextBuilderwith required parameters.- Parameters:
path – The root path of the table.
Constructs a
ReadContextBuilderfor a format table that is already loaded: the only way to read one whose schema lives in a metastore rather than under its location, such as a table a REST catalog serves.The table carries what such a location does not say, so
SetTableSchema(),WithFileSystem(),WithFileSystemSchemeToIdentifierMap()and a branch are refused here rather than ignored.- Parameters:
table – The format table to read, as
Catalog::GetFormatTable()hands it back.
-
~ReadContextBuilder()#
-
ReadContextBuilder(ReadContextBuilder&&) noexcept#
-
ReadContextBuilder &operator=(ReadContextBuilder&&) noexcept#
-
ReadContextBuilder &SetReadFieldNames(const std::vector<std::string> &read_field_names)#
Set the schema fields to read from the table.
If not set, all fields from the table schema will be read. This is useful for projection pushdown to reduce I/O and improve performance by reading only the required columns.
Note
Currently supports top-level field selection. For nested field selection use SetReadSchema(std::unique_ptr<ArrowSchema>) instead.
- Parameters:
read_field_names – Vector of field names to read from the table.
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &SetReadFieldIds(const std::vector<int32_t> &read_field_ids)#
Set the schema fields to read from the table.
If not set, all fields from the table schema will be read. This is useful for projection pushdown to reduce I/O and improve performance by reading only the required columns.
Note
Currently supports top-level field selection.
Note
SetReadFieldIds() and SetReadFieldNames() are mutually exclusive. Calling both will ignore the read schema set by SetReadFieldNames().
- Parameters:
read_field_ids – Vector of field ids to read from the table.
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &SetReadSchema(std::unique_ptr<ArrowSchema> read_schema)#
Set the read Arrow Schema for nested column pruning.
The read schema is an Arrow C Data Interface schema where STRUCT types may contain only a subset of the original sub-fields, enabling nested column pruning to reduce I/O. Field matching is based on field name: the system looks up each field by name in the table schema and rebuilds the aligned schema using the table schema’s type and metadata. Metadata propagation from the user-provided schema is whitelist-based: currently only “paimon.map.selected-keys” is preserved and merged into the final aligned schema.
To prune map entries by key, attach metadata “paimon.map.selected-keys” to the target map field in read schema. The value is a comma-separated key list, for example: “k1,k2”. Only map fields with string key type (Arrow utf8) are supported.
Attaching this metadata to a MAP field only filters the returned MAP after reading. To push down selected keys from a shared-shredding MAP and return them as STRUCT children, build the field with
MapSharedShreddingAccessBuilder. To read selected paths from a VARIANT field, build the field withVariantAccessBuilder.Example:
auto map_field = arrow::field("m", arrow::map(arrow::utf8(), arrow::int32())); auto map_meta = arrow::KeyValueMetadata::Make( {"paimon.map.selected-keys"}, {"k1,k2"}); auto projected_schema = arrow::schema({ arrow::field("id", arrow::int64()), map_field->WithMetadata(map_meta), }); auto c_schema = std::make_unique<ArrowSchema>(); arrow::ExportSchema(*projected_schema, c_schema.get()); ReadContextBuilder builder("/path/to/table"); builder.SetReadSchema(std::move(c_schema));
Note
Priority: read_schema > read_field_ids > read_field_names. When set, read_field_ids and read_field_names are ignored.
- Parameters:
read_schema – Arrow C Schema. Ownership of schema resources is transferred to the built ReadContext.
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &SetOptions(const std::map<std::string, std::string> &options)#
Set a configuration options map to set some option entries which are not defined in the table schema or whose values you want to overwrite.
Note
The options map will clear the options added by
AddOption()before.- Parameters:
options – The configuration options map.
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &AddOption(const std::string &key, const std::string &value)#
Add a single configuration option which is not defined in the table schema or whose value you want to overwrite.
If you want to add multiple options, call
AddOption()multiple times or useSetOptions()instead.- Parameters:
key – The option key.
value – The option value.
- Returns:
Reference to this builder for method chaining.
Set a predicate for filtering data during reading.
The predicate is used for both partition pruning and data filtering. It can significantly improve performance by reducing the amount of data that needs to be read and processed.
- Parameters:
predicate – Shared pointer to the predicate for data filtering.
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &EnablePredicateFilter(bool enabled)#
Whether to perform precise filtering according to predicates for data read from format reader.
- Parameters:
enabled – Whether to enable precise filtering (default: false)
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &EnablePrefetch(bool enabled)#
Enable or disable prefetching of data batches from individual files.
When enabled, the reader will prefetch multiple batches in parallel to improve throughput by overlapping I/O with computation. This is particularly beneficial for high-latency storage systems.
- Parameters:
enabled – Whether to enable prefetching (default: false)
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &EnableLateMaterializing(bool enabled)#
Enable or disable late materialization (probe/payload two-phase reads).
When enabled, each parallel reader under the prefetch layer performs a probe read of predicate columns first and only materializes payload columns for matched rows.
Note
Without a pushed-down predicate the late-materializing reader degrades to a plain passthrough.
- Parameters:
enabled – Whether to enable late materialization (default: false)
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &SetReadAheadCacheEnabled(bool enabled)#
Enable or disable the read-ahead cache for read operations.
A read-ahead cache is used to prebuffer data ranges before they are needed, which can improve read performance by reducing redundant I/O operations.
- Parameters:
enabled – Whether to enable the read-ahead cache (default: true)
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &WithCacheConfig(const CacheConfig &config)#
Set the cache configuration for prefetch read operations.
- Parameters:
config – The cache configuration to use.
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &SetWarmupLevel(WarmupLevel level)#
Set how far the reader prepares the next file before it is read.
Warmup overlaps remote-storage latency with the read of the current file. Higher levels hide more latency but use more memory, and may warm files that a query never reads (for example when a LIMIT stops the scan early).
See also
Note
WarmupLevel::RAW warms the read-ahead cache, so it has no effect and behaves like WarmupLevel::NONE when the cache is off (see SetReadAheadCacheEnabled()). WarmupLevel::DECODED does not depend on the cache.
- Parameters:
level – The warmup level to use (default: WarmupLevel::RAW).
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &SetPrefetchBatchCount(uint32_t batch_count)#
Set the total number of batches to prefetch across all files.
This controls the memory usage and parallelism of the prefetching mechanism. Higher values can improve throughput but consume more memory.
- Parameters:
batch_count – Total number of batches to prefetch (default: 600)
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &SetPrefetchMaxParallelNum(uint32_t parallel_num)#
Set the maximum number of parallel prefetch operations.
This limits the number of concurrent I/O operations to prevent overwhelming the storage system or consuming excessive system resources.
- Parameters:
parallel_num – Maximum parallel prefetch operations (default: 3)
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &EnableMultiThreadRowToBatch(bool enabled)#
Enable or disable multi-threaded row-to-batch conversion in merge-on-read scenarios.
When enabled, multiple threads are used to convert row data to batch format during merge operations, which can improve performance for CPU-intensive merge operations.
- Parameters:
enabled – Whether to enable multi-threaded conversion (default: false)
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &SetRowToBatchThreadNumber(uint32_t thread_number)#
Set the number of threads for row-to-batch conversion in merge-on-read scenarios.
This controls the parallelism of row-to-batch conversion during merge operations. Higher values can improve performance but may affect result ordering.
Note
If thread_number > 1, Arrow batches from the reader may not be in primary key order.
- Parameters:
thread_number – Number of conversion threads (default: 1)
- Returns:
Reference to this builder for method chaining.
Set custom memory pool for memory management.
Note
If not set, the default system memory pool will be used.
- Parameters:
memory_pool – The memory pool to use.
- Returns:
Reference to this builder for method chaining.
Set custom executor for task execution.
Note
If not set, the default system executor will be used.
- Parameters:
executor – The executor to use.
- Returns:
Reference to this builder for method chaining.
Sets the shared context used to resolve process-local real-time split tickets.
- Parameters:
realtime_context – Context shared with the scan that created the real-time splits.
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &SetTableSchema(const std::string &table_schema)#
Set the table schema as a string to avoid schema loading I/O operations.
This optimization allows the reader to use a pre-loaded schema instead of reading it from the table metadata, which can improve performance especially in scenarios with many small read operations.
Note
The user must ensure that the schema string is valid and matches the table.
Note
If not set, the schema will be loaded from the table path.
- Parameters:
table_schema – String representation of the table schema.
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &WithBranch(const std::string &branch)#
Set the specific branch to read from in a versioned table.
Paimon supports branching for data versioning and time travel queries. This method allows reading from a specific branch instead of the main branch.
The
branchoption names the branch too, as it does for a scan of one branch; naming two different branches is refused rather than silently resolved.Note
Default branch is “main” if not specified. An empty name is the main branch.
- Parameters:
branch – Name of the branch to read from.
- Returns:
Reference to this builder for method chaining.
-
ReadContextBuilder &WithFileSystemSchemeToIdentifierMap(const std::map<std::string, std::string> &fs_scheme_to_identifier_map)#
Sets a mapping from URI schemes (e.g., “file”, “oss”) to registered file system identifiers.
This allows selecting different pre-registered file system implementations based on the URI scheme at runtime.
Note
This method is intended for environments where multiple file systems are pre-registered.
The specified identifiers must correspond to file systems that have been registered at compile time or initialization.
Cannot be used together with
WithFileSystem().If not set, use default file system (configured in
Options::FILE_SYSTEM). Example: builder.WithFileSystemSchemeToIdentifierMap({{“oss”, “jindo”}, {“file”, “local”}});
- Parameters:
fs_scheme_to_identifier_map – Map from URI scheme (like “oss”) to the corresponding file system identifier.
- Returns:
Reference to this builder for method chaining.
Sets a custom file system instance to be used for all file operations in this read context.
This bypasses the global file system registry and uses the provided implementation directly.
Note
If not set, use default file system (configured in
Options::FILE_SYSTEM)- Parameters:
file_system – The file system to use.
- Returns:
Reference to this builder for method chaining.
Reads a native table through its own file system - including the per-table temporary credentials a catalog that issues them hands out through
Catalog::GetTableFileSystem.This is a shorthand for
WithFileSystem(catalog->GetTableFileSystem(identifier)); an explicitWithFileSystem()takes precedence, so the catalog is not asked.- Parameters:
catalog – Non-null catalog, read when
Finish()builds the context.identifier – The native table to read.
- Returns:
Reference to this builder for method chaining.
Inject a cache for read operations.
Passing nullptr disables cache.
- Returns:
Reference to this builder for method chaining.
-
Result<std::unique_ptr<ReadContext>> Finish()#
Build and return a
ReadContextinstance with input validation.- Returns:
Result containing the constructed
ReadContextor an error status.
-
explicit ReadContextBuilder(const std::string &path)#
-
enum class paimon::WarmupLevel#
Controls how far a reader prepares the next file before that file is actually read.
Warmup overlaps remote-storage latency with the read of the current file. Each level takes the next file one step further along the read pipeline:
RAWcovers the remote fetch,DECODEDadds the decode on top of it. Higher levels hide more latency, but commit more memory and background I/O to files that a query may end up never reading (for example when a LIMIT stops the scan early). Callers can trade latency against memory by picking a level.Values:
-
enumerator NONE#
Do not warm up.
The next file’s I/O starts only when it is actually read. This is the behavior from before warmup existed and uses no extra memory or background threads.
-
enumerator RAW#
Fetch only the next file’s raw, still-compressed bytes into memory, and leave the decoder alone.
Overlaps the remote fetch while keeping memory lower than
DECODED, because no decoded batches are materialized ahead of the read. It fetches through the read-ahead cache, so it falls back toNONEwhen that cache is disabled. This is the default, as it hides the remote fetch without committing memory to decoded batches.
-
enumerator DECODED#
Fetch the raw bytes and start the background decode loop as well, so decoded batches are ready before the file is read.
Hides the most latency but uses the most memory.
-
enumerator NONE#
-
class ReadContext#
ReadContextis some configuration for read operations.Please do not use this class directly, use
ReadContextBuilderto build aReadContextwhich has input validation.See also
Public Functions
-
~ReadContext()#
-
inline const std::string &GetPath() const#
-
inline const std::string &GetBranch() const#
-
inline const std::map<std::string, std::string> &GetFileSystemSchemeToIdentifierMap() const#
-
inline const std::map<std::string, std::string> &GetOptions() const#
-
inline const std::vector<std::string> &GetReadFieldNames() const#
-
inline const std::vector<int32_t> &GetReadFieldIds() const#
-
inline bool EnablePredicateFilter() const#
-
inline bool EnablePrefetch() const#
-
inline bool EnableLateMaterializing() const#
-
inline uint32_t GetPrefetchBatchCount() const#
-
inline uint32_t GetPrefetchMaxParallelNum() const#
-
inline bool EnableMultiThreadRowToBatch() const#
-
inline uint32_t GetRowToBatchThreadNumber() const#
-
inline const std::optional<std::string> &GetSpecificTableSchema() const#
-
inline std::shared_ptr<MemoryPool> GetMemoryPool() const#
-
inline std::shared_ptr<FileSystem> GetSpecificFileSystem() const#
-
inline std::shared_ptr<RealtimeContext> GetRealtimeContext() const#
Returns the context used to resolve process-local real-time split tickets.
-
inline bool ReadAheadCacheEnabled() const#
-
inline const CacheConfig &GetCacheConfig() const#
-
inline std::shared_ptr<Cache> GetCache() const#
-
inline const std::shared_ptr<FormatTable> &GetFormatTable() const#
The format table this context was built from, or null when it names a table path and the schema under that path says what kind of table it is.
-
inline WarmupLevel GetWarmupLevel() const#
-
inline bool HasReadSchema() const#
Whether a read schema (C ArrowSchema) for nested column pruning was provided.
-
inline ArrowSchema *GetReadSchema()#
Get the read schema as a mutable C ArrowSchema pointer.
ImportSchema will consume (release) the schema content.
-
void SetReadSchema(std::unique_ptr<ArrowSchema> schema)#
Set the read schema from a C ArrowSchema unique_ptr and take ownership of schema resources (released via ArrowSchema::release in destructor).
Called internally by ReadContextBuilder.
-
~ReadContext()#
-
class BatchReader#
A batch reader that supports reading batch data into an arrow array.
Subclassed by paimon::FileBatchReader
Public Types
-
using ReadBatch = std::pair<std::unique_ptr<ArrowArray>, std::unique_ptr<ArrowSchema>>#
Public Functions
-
virtual ~BatchReader() = default#
-
virtual Result<ReadBatch> NextBatch() = 0#
Retrieves the next batch of data.
If EOF is reached, returns an OK status with a nullptr array. Returns an error status only for critical failures (e.g., IO errors). Once an error is returned, this method must not be retried, as it will repeatedly return the same error code.
Warning
A non-EOF ArrowArray and all its nested child arrays must have offset 0 to avoid potential issues during conversion through the Arrow C Data Interface.
Warning
A returned ArrowArray must retain every allocator and plugin resource needed by its release callback, so it remains releasable after this reader is destroyed.
Warning
Consumers must treat the returned ArrowArray and ArrowSchema as one complete Arrow C Data Interface ownership unit. Moving or retaining an individual child ArrowArray without its root array is unsupported because resource lifetimes are retained by the root array’s release chain.
- Returns:
A result containing a
ReadBatch, which consists of a unique pointer toArrowArrayand a unique pointer toArrowSchema. Returned array contains a_VALUE_KINDfield (the first field) to indicate the row kind of each row. Deleted or index-filtered rows are removed.
-
virtual Result<ReadBatchWithBitmap> NextBatchWithBitmap()#
Retrieves the next batch of data.
If EOF is reached, returns an OK status with a nullptr array. Returns an error status only for critical failures (e.g., IO errors). Once an error is returned, this method must not be retried, as it will repeatedly return the same error code.
Warning
A non-EOF ArrowArray and all its nested child arrays must have offset 0 to avoid potential issues during conversion through the Arrow C Data Interface.
Warning
A returned ArrowArray must retain every allocator and plugin resource needed by its release callback, so it remains releasable after this reader is destroyed.
Warning
Consumers must treat the returned ArrowArray and ArrowSchema as one complete Arrow C Data Interface ownership unit. Moving or retaining an individual child ArrowArray without its root array is unsupported because resource lifetimes are retained by the root array’s release chain.
- Returns:
A result containing a
ReadBatchand a valid bitmap.ReadBatchconsists of a unique pointer toArrowArrayand a unique pointer toArrowSchema. Returned array contains a _VALUE_KIND field (the first field) to indicate the row kind of each row. Deleted or index-filtered records maybe maintained inReadBatch, while bitmap indicates valid row id. If deletion vector or index are enabled, this function is more efficient thanNextBatch(). The default implementation callsNextBatch()and adds all rows to valid bitmap. Noted that the returned bitmap has at least one valid row id.
-
virtual std::shared_ptr<Metrics> GetReaderMetrics() const = 0#
Retrieves the reader’s metrics.
Note that calling this method frequently may incur significant performance overhead.
- Returns:
A shared pointer to the
Metricsobject.
-
virtual void Close() = 0#
Closes the
BatchReader, releasing any associated resources.After calling this method, further calls to
NextBatch()is undefined and should be avoided.
Public Static Functions
-
static bool IsEofBatch(const ReadBatch &batch)#
Determine whether a
ReadBatchorReadBatchWithBitmapis eof batch, if return true, all the data has been returned.
-
static bool IsEofBatch(const ReadBatchWithBitmap &batch_with_bitmap)#
-
static ReadBatchWithBitmap MakeEofBatchWithBitmap()#
-
using ReadBatch = std::pair<std::unique_ptr<ArrowArray>, std::unique_ptr<ArrowSchema>>#