alibaba / alibaba/paimon-cpp

[Feature] Support pluggable real-time writes and memory/disk union reads

Open
#445 0 comments 0 reactions 0 assignees View on GitHub
enhancement
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

Open the contributing 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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.