[Feature] Support pluggable real-time writes and memory/disk union reads
- Dominant language
- C++
- Stars
- 131
- Forks
- 49
- PR merge metrics
- No merged PRs in 30d
Description
### Search before asking
- [x] I searched in the [issues](https://github.com/alibaba/paimon-cpp/issues) and found nothing similar.
### Motivation
Paimon data is normally queryable only after it is written to data files and committed into a snapshot. Some workloads need to query data still held by the writer while continuing ingestion during `PrepareCommit`.
We propose a pluggable real-time layer for append and primary-key tables while preserving Paimon's existing routing, file format, manifest, snapshot, and commit semantics.
### Solution
Introduce one `MemIndexer` for each routed partition-bucket and manage its data as building, sealed, and reclaimable segments.
- Paimon assigns monotonically increasing sequence numbers within each partition-bucket.
- `PrepareCommit` seals the current segment and immediately opens a new writable segment.
- Paimon writes sealed data through inner existing rolling writers and produces standard commit messages.
- Each snapshot persists the committed sequence watermark of every partition-bucket.
- A query reads the committed snapshot plus memory rows whose sequence is above the committed watermark.
- After a new snapshot covers a sealed segment, new queries switch to disk and the segment is reclaimed after existing readers release it.
### Plugin API
The following is an API sketch for discussion:
```cpp
struct RealtimeWriteBatch {
std::unique_ptr batch;
Range sequence_range;
};
class RealtimeSegmentHandle {
public:
virtual Range GetSequenceRange() const = 0;
};
class MemIndexer {
public:
virtual Status Write(RealtimeWriteBatch&& batch) = 0;
virtual Result>>
SealForCommit() = 0;
// MemReadRequest selects the visible sequence range and memory segments.
virtual Result> AcquireReadView(
const MemReadRequest& request) = 0;
// MemReadView pins that selection as a stable, reference-protected view.
// MemQueryContext supplies the projection and predicate when creating readers.
virtual Result>>
CreateQueryReaders(const std::shared_ptr& view,
const MemQueryContext& context) = 0;
virtual Result>>
CreateCommitReaders(
const std::shared_ptr& segment) = 0;
virtual Status Reclaim(
const std::shared_ptr& segment) = 0;
};
class MemIndexerFactory {
public:
virtual Result> Create(
const MemIndexerOptions& options) = 0;
};
```
Paimon performs schema validation, partition-bucket routing, sequence assignment, table-specific merging, and file writing. The plugin manages memory or spill storage and creates query and commit readers.
An opaque `RealtimeContext` owns the partition-bucket to `MemIndexer` mapping and is shared by write and scan operations:
```cpp
auto realtime = RealtimeContext::Create(mem_indexer_factory);
WriteContextBuilder(...).WithRealtimeContext(realtime);
ScanContextBuilder(...).WithRealtimeContext(realtime);
```
### Union Read
The existing reader interface remains unchanged:
```cpp
virtual Result> CreateReader(
const std::shared_ptr& split) = 0;
```
When real-time reading is enabled, `TableScan::CreatePlan()` atomically captures:
```text
committed snapshot S
disk watermark D for each partition-bucket
mem upper watermark U
MemIndexer and MemReadView references
```
It first pins the memory views and then creates disk splits for snapshot `S`. Disk splits and memory views are grouped by partition-bucket into internal `RealtimeSplit` objects:
```text
RealtimeSplit
disk split(s)
disk watermark D
mem upper watermark U
MemIndexer reference
MemReadView reference
```
A `RealtimeSplit` may contain only disk data or only memory data.
`TableRead::CreateReader(split)` handles both split types:
```text
DataSplit
-> disk BatchReader
RealtimeSplit
-> disk reader(s)
-> MemIndexer query reader(s)
-> append concat or primary-key merge
```
For append tables:
```text
disk readers + mem readers -> ConcatBatchReader
```
For primary-key tables:
```text
disk KeyValue readers + sorted mem runs
-> existing sorted merge
-> existing MergeFunction
-> BatchReader
```
The plugin does not implement primary-key comparison, deduplication, deletion, partial-update, or aggregation semantics. These remain in Paimon's existing merge pipeline.
A `RealtimeSplit` and its resulting reader retain the pinned `MemReadView`, so referenced segments cannot be reclaimed while an older query is still running. The initial implementation is process-local; serializable or remote real-time splits can be considered separately.
### Anything else?
The implementation can be incremental:
The initial scope assumes fixed buckets, stable schemas, one writer per partition-bucket, and concurrent readers.
1. Add the plugin, segment lifecycle, sequence progress.
2. Add arrow-based default plugin.
3. Support append table.
4. Support pk table (MOR).
5. Add optional predicate indexes and precise key lookup optimizations.
6. Support dv mode.
### Are you willing to submit a PR?
- [x] I'm willing to submit a PR!
Contributor guide
Research direction
Start by tracing TableScan::CreatePlan() and TableRead::CreateReader(), then review the WriteContextBuilder and ScanContextBuilder integration points and the existing sorted merge pipeline. Done means adding the incremental MemIndexer lifecycle, union reads, and append and primary-key support while preserving existing commit and merge semantics.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- cpp
- Domain
- databases
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Quiet
- Clarity
- Mostly clear
- Newbie friendliness
- 35/100