[Feature] Prefetch the late-materialization payload ranges instead of reading them on demand
- 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
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