apache / apache/paimon-cpp

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

Fechada
#341 0 comentários 0 reações 0 responsáveis Ver no GitHub
enhancement
Linguagem predominante
C++
Estrelas
65
Forks
25
Merge médio
2d 12h
PRs com merge (30d)
80

Descrição

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

Guia de contribuição

Abrir o guia de contribuição

Direção de pesquisa

Comece rastreando PrefetchFileBatchReader::PreBufferRange(), ReadAheadCache::Init(), Read(), Reset() e Close(), depois inspecione as APIs públicas em include/paimon/. Acompanhe como LateMaterializingFileBatchReader determina os intervalos de payload e como as métricas de cache existentes são registradas. O trabalho estará concluído quando os intervalos tardios forem registrados e aquecidos com segurança entre as rodadas, os registros obsoletos forem descartados, as regras de concorrência e sobreposição forem mantidas e as novas métricas contabilizarem os bytes registrados e descartados.

Escrita pelo modelo de indexação a partir do texto da issue.

Avaliação

Stack de tecnologia
cpp
Domínio
data-engineering, performance
Tipo de issue
Funcionalidade
Dificuldade
5/5
Tempo estimado
Mais de uma semana
Status de atividade
Ativa
Clareza
Claramente especificada
Facilidade para iniciantes
38/100

Receba novas issues na sua caixa de entrada

Um resumo curto de issues do GitHub para quem está começando.