Skip to content

[RFC] Partition statistics for catalog-managed format tables #9123

Description

@sundapeng

This is a design discussion, not a patch. The first four PRs of the stack in section 7 are open (#9119, #9120, #9121, #9122); the rest are held until the contract and the protocol below get a ruling.

1. Problem

A catalog-managed format table registers its partitions with the catalog, and the catalog returns
those registrations to every engine that plans a query against it. Today it returns no statistics
for them — and worse than no, it returns numbers that look real.

Two things are wrong, and they are independent:

(a) The contract does not say what a negative statistic means. PartitionStatistics is read on
two planes. On the delta plane — what a commit changed — a negative value is a decrement the server
adds to what it holds; that is what the current javadoc describes. On the observation plane — what
listPartitions returns for a partition as it stands — a negative value can only mean "nobody ever
reported this". The javadoc covers only the first, so consumers have been reading the second as if it
were the first, or as if unknown were zero.

Treating unknown as zero is not a cosmetic error. A planner that reads zero rows plans against an
empty partition that may hold a billion. An aggregate that answers COUNT(*) from a zero row count
returns zero for a partition full of data, without touching a file and without any error — this has
happened, in apache/paimon-rust #624.

(b) There is no channel for the numbers at all. A format table has no snapshot, so
commitSnapshot — the channel a table snapshot uses to report statistics — does not exist for it.
Meanwhile the write path already has the numbers and throws them away: the rolling writer counts
every row it writes, the single-file writer knows the byte length of the file it closed, and
prepareCommit wraps each committer in a TwoPhaseCommitMessage that carries neither.

So the numbers exist, cost nothing to collect, and have nowhere to go.

2. Non-goals

  • Column statistics. This is partition level only.
  • Table-level statistics. A format table has no snapshot to carry them.
  • Changing the delta plane. Negative deltas stay exactly as they are.
  • Deciding which partitions exist. Statistics never create or remove a partition row.
  • Deriving statistics from a diff against object storage. The number of live files a table
    semantically holds and the number of objects a bucket physically holds are different quantities;
    their difference must not drive deletion.

3. The contract

Plane Where it comes from Negative means Zero means
Delta a commit reporting what it changed a decrement to apply no change
Observation listPartitions never measured (UNKNOWN) an exact zero

Three rules follow, and all three matter:

  1. Unknown is per field. A reporter that can measure the file count but not the row count leaves
    the row count unknown and fills the rest. It must not have to choose between reporting a guess and
    reporting nothing.
  2. Unknown is not zero. An unreported partition is not an empty partition. Anything that filters,
    sums or timestamps a statistic has to distinguish them.
  3. The fields stay primitive. Boxing them to express unknown as null would break a @Public
    class, and the encoding above needs no new type.

PartitionStatistics.UNKNOWN names the canonical value and isKnown(long) tests it, so callers stop
comparing against -1 at each site.

4. Reporting semantics

Which mode a report uses follows from what the reporter did, not from a preference:

Write Mode Target partitions Why it is exact
append ADD the ones written the writer saw only its own files, so it can only report an increment
dynamic overwrite SET the ones written the write replaced everything those partitions held, so what it wrote is the total
static prefix overwrite SET the whole cleared subtree ⊋ the ones written same, plus partitions this commit wrote nothing to

Static prefix overwrite is why a pure increment cannot express this. Clearing a prefix empties
every partition beneath it, including ones the commit writes nothing to. Their old data is gone and
no increment says so. They report zero — an exact zero, they really are empty — and stay
registered
. The directories emptied that way come out of the deletion listing, which already had to
walk them, so those numbers cost no extra IO.

Per field, per mode:

Field ADD SET
recordCount / fileSizeInBytes / fileCount summed replaced
lastFileCreationTime max of stored and reported replaced
any field reported unknown skipped skipped

A timestamp is not a quantity: summing two epoch-millisecond values produces a meaningless number, so
ADD takes the later of the two. SET may move it backwards, which is correct — the reporter saw the
whole partition, so the newest file it found is the newest there is.

4.1 ADD is not idempotent, and this RFC does not pretend otherwise

A redelivered ADD counts twice, and nothing in the report lets the server tell a redelivery from a
second genuine increment. Three sources of drift, none of which this closes:

Source Note
a redelivered request the RPC succeeded but the response was lost
a writer that is not Paimon files appear that no commit reported
a file deleted out of band files disappear that no commit reported

Convergence is a later full report over the same partition — which is what the two management
commands in §5 are for. A deduplication token in front of every commit would need server-side state
with its own lifecycle, and would cost more than the drift does. This is a trade-off, not an
oversight
, and it is the open question most worth a maintainer's opinion (§8).

What this RFC does close is the case where the client causes the duplication by itself: a POST
carrying a non-empty ADD report declares itself unsafe to replay, so the retry that a 429 or a 503
triggers cannot count the same increment twice. See §6.

5. Where the numbers come from

Source Row count File count / bytes / time Extra IO Covers
commit exact, already counted exact, already known none only what Paimon wrote
MSCK REPAIR TABLE — (listing only) exact O(files) listing everything, whoever wrote it
ANALYZE TABLE exact where a footer exists exact listing + footers everything

The commit is the only place a row count is free. A rescan pays a file footer for it, and only
formats that carry one can give it at all: a CSV, TEXT or JSON partition keeps an unknown row count
rather than a guessed one. One unreadable footer makes the whole partition's row count unknown rather
than short — a sum missing a file, reported as exact, is worse than no number.

Both management commands measure only partitions that are already registered, and both replace
rather than accumulate, so running either twice is running it once. Measuring must never register a
partition the command was not asked to.

Both are opt-in, for the same reason: measuring changes what the command costs. A plain repair
lists partition directories; measuring lists the files inside every one of them. That is a different
order of magnitude on a table with many partitions, and a command should not silently become that.

A listing failure aborts the whole collection rather than reporting what it managed to see: a
truncated listing is indistinguishable from a partition that lost files.

6. Protocol

The registration request carries the statistics. Two optional fields on the existing
create-partitions request — the statistics list and the mode — and nothing else.

Why not a separate endpoint. Registration and statistics then land in one request and one
server-side transaction. "The statistics failed but the partition registered" stops being a state
anyone has to handle, and no extra round trip is paid. A separate endpoint buys exactly one thing:
the ability to report statistics without registering — which is the state worth not having.

Compatibility runs both ways and neither direction errors.

Behaviour
new client, old server the field is ignored; statistics stay unknown; the client logs once that the server did not accept them
old client, new server no statistics arrive; the columns stay unknown
a caller that reports nothing sends exactly the request it sends today, so the server sees no change in shape

The asymmetry to be explicit about: an observation that is missing is something consumers can
handle; an observation that is wrong is not.
Silently dropping statistics is an acceptable
degradation only because §3 makes unknown expressible. Without §3 the same silent drop would leave
a fabricated zero behind, and would have to fail closed instead. §3 is a prerequisite for §6, not a
parallel change.

6.1 Retry

RESTRequest gains isRetrySafe(), defaulting to true so every existing request keeps the 429/503
retry it has today — that retry is the only defence against a rate limiter or a restarting node, and
nearly every request Paimon sends over POST is idempotent by content even though POST is not
idempotent by method.

A request answering false is sent exactly once and the failure reaches the caller. The mark travels
in the client context rather than in the request, so it never reaches the wire and survives whatever
the exec chain does to the request object; the getter is annotated so it stays out of the serialized
body as well.

7. PR stack

One diff would put a contract clarification, an internal data-flow change, a wire change, a
behaviour change and two management commands in front of one reviewer. Eight PRs instead:

# Title Depends on Reviewable on its own because
S0 [core] Spell out what a negative partition statistic means contract + named constant, zero behaviour change
S1 [core] Carry the row count and byte size a format table writer already counted S0 paimon-core internal, nothing reads the numbers yet
S2 [rest] Do not replay a POST the server cannot absorb twice opt-in per request, default keeps today's behaviour
S3 [core][rest] Report partition statistics when registering partitions S0, S2 new API is all defaults and overloads
S4 [core] Report what a format table commit wrote to the catalog S0, S1, S3 the only behaviour change, behind a flag, off by default
S5 [core][spark] Measure format table partitions in MSCK REPAIR TABLE S0, S3 off by default
S6 [spark] Support ANALYZE TABLE … COMPUTE STATISTICS on catalog-managed format tables S5 a different command with a different cost model
S7 [spark] Cover the partition operations left over tests plus one error message
S7 ──────────────────────────────────────────► independent
S0 ──┬──► S1 ──┐
     │         ├──► S4
S2 ──┴──► S3 ──┴──► S5 ──► S6

Two orderings are not negotiable:

  1. S2 before S3. A non-empty ADD report is not retry-safe. With S3 in and S2 out, an automatic
    retry reports the same increment twice and the server cannot tell.
  2. S4 alone. It is the only PR that changes what an existing write does.

If S3–S6 are rejected, S0 still stands on its own and is worth having: the ambiguity it removes has
already produced one silent wrong answer.

8. Open questions for maintainers

  1. Is "ADD may drift, converge with a later SET" acceptable, or should the protocol carry an
    idempotency token from the start? This decides whether S3/S4 are shaped as they are. §4.1 argues
    for accepting the drift; the counter-argument is that a wrong number in a catalog is a wrong number
    regardless of how it got there.
  2. Should the client be able to tell that a server ignored its statistics, beyond logging once?
    A capability flag in the response is the obvious answer, and the obvious cost is one more thing to
    version.
  3. Should lastFileCreationTime come from the commit clock or from each file's mtime? The clock
    is free; the mtimes cost one request per file and give a coarser answer for a partition written
    over a long window. Currently the clock.
  4. A partition emptied by a static prefix overwrite keeps its old lastFileCreationTime (SET
    only moves it forward under ADD, and there is no file to date under SET). Known conservative
    choice — worth a second opinion.

9. Verification matrix

Synthetic values throughout; a partition written by a single commit of 3 rows into one file.

Scenario Expected
registered, never reported every field unknown, per field
append 3 rows, then 2 more ADD → 5 rows, 2 files
dynamic overwrite writing 4 rows over those 5 SET → 4, not 9
static prefix overwrite, sibling partition cleared but not rewritten sibling → exact 0, still registered
report all zeros partition still returned by listPartitions
CSV table, rescan file count / bytes / time filled, row count still unknown
one unreadable footer in a partition that partition's row count unknown, not short
listing fails mid-rescan whole collection aborts, nothing written
ADD reported, response lost, request retried sent once; the failure reaches the caller
repeat the same SET idempotent
DROP PARTITION row gone; statistics go with it; no separate correction needed

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions