apache / apache/paimon-cpp

[Feature] Prefetch the late-materialization payload ranges instead of reading them on demand

Abierto
#341 0 comentarios 0 reacciones 0 asignados 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

Late materialization reads a data file in two passes: a probe pass over the predicate fields, then a payload pass over the remaining fields for the matched rows only. The shared read-ahead cache is fed once per read-range generation through `PrefetchFileBatchReader::PreBufferRange()`, before any read starts. At that point the payload pass cannot know which pages hold the matched rows — that depends on the probe result — so `PreBufferRange()` only reports the probe ranges (there was an explicit TODO for exactly this). The payload pass is therefore never prefetched: every payload read misses the cache and waits for its own underlying IO, serialized against the decode, on the pass that touches the wide columns.

## Solution

Let a reader report byte ranges that only become known after reading has started, and let the shared cache register them mid-read.

- New `ReadAheadCache::AddRanges(ranges, expected_round)` registers ranges into an already-initialized cache and is safe to call repeatedly and concurrently with `Read()`. It merges the new ranges into the disjoint, offset-ordered pending list, registering only the parts no registered range covers and dropping the overlap (the round that registered it is already fetching those bytes), then rebuilds the per-range cached flags so an already-fetched range is not fetched twice. The registered part is cut at a new `CacheConfig` knob `late_range_size_limit` (default 8 MiB, smaller than the 32 MiB `range_size_limit`) so a large pass is fetched by several concurrent requests rather than one long one; a new `Warmup(from_offset)` starts fetching from the first newly-registered range instead of from the head.
- A registration round bounds the lifetime: every `Init()` opens a round identified by `RegistrationRound()`, and `AddRanges()` drops everything when `expected_round` is not the open round, so a pass that outlived its generation — the cache was reset for a new read-range generation, or released by `Close()` — registers nothing instead of prefetching bytes nobody reads. The round counter is monotonic across `Reset()` so a stale round is never mistaken for a new one.
- New `PrefetchFileBatchReader::PreBufferSink` and `SetPreBufferSink()`: `PrefetchFileBatchReaderImpl` installs a sink on each sub-reader that tags the reported ranges with the current round, calls `AddRanges`, and warms up from the first new range. `LateMaterializingFileBatchReader` reports the payload ranges through the sink once the probe pass has refined the inner reader's target pages, and surfaces a failure to compute them (they come from the file metadata) rather than swallowing it.
- New metrics `read-ahead-cache.late.registered` / `.registered-bytes` / `.dropped` / `.dropped-bytes`, counted after coalescing and splitting, so `registered-bytes` and `dropped-bytes` together account for every reported byte.

## Anything else?

Adds public API under `include/paimon/`: `PrefetchFileBatchReader::PreBufferSink` / `SetPreBufferSink()` and `CacheConfig::GetLateRangeSizeLimit()` / `SetLateRangeSizeLimit()`. `ReadAheadCache::AddRanges` / `RegistrationRound` / `Warmup(offset)` and the new counter names live in the internal header. No storage format or protocol change.

## 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

Línea de trabajo

Comienza siguiendo PrefetchFileBatchReader::PreBufferRange(), ReadAheadCache::Init(), Read(), Reset() y Close(), y después inspecciona las APIs públicas bajo include/paimon/. Sigue cómo LateMaterializingFileBatchReader determina los rangos de payload y cómo se registran las métricas de caché existentes. La tarea estará terminada cuando los rangos tardíos se registren y se calienten de forma segura entre rondas, se descarten los registros obsoletos, se cumplan las reglas de concurrencia y solapamiento, y las nuevas métricas contabilicen los bytes registrados y descartados.

Escrito por el modelo de indexación a partir del texto del issue.

Evaluación

Stack tecnológico
cpp
Área
data-engineering, performance
Tipo de issue
Nueva funcionalidad
Dificultad
5/5
Tiempo estimado
Más de una semana
Estado de actividad
Activo
Claridad
Bien especificado
Aptitud para principiantes
38/100

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.