apache / apache/paimon-cpp

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

Abierto
#158 0 comentarios 0 reacciones 1 asignado Reclamado por @lxy-9602 Ver en GitHub
enhancement
Lenguaje dominante
C++
Estrellas
65
Forks
25
Merge medio
2 d 12 h
PR fusionados (30 d)
80

Descripción

### Search before asking

- [x] I searched in the [issues](https://github.com/apache/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!

Guía de contribución

Abrir la guía de contribución

Evaluación

Este issue todavía no se ha evaluado.

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.