Skip to content

[core] Support adaptive MAP and Variant shredding across rolling files#8827

Open
lxy-9602 wants to merge 4 commits into
apache:masterfrom
lxy-9602:map-variant-restore
Open

[core] Support adaptive MAP and Variant shredding across rolling files#8827
lxy-9602 wants to merge 4 commits into
apache:masterfrom
lxy-9602:map-variant-restore

Conversation

@lxy-9602

@lxy-9602 lxy-9602 commented Jul 23, 2026

Copy link
Copy Markdown
Contributor

Purpose

This is a subtask of MAP shared-shredding support.

This PR stabilizes shredding layouts across files created by the same rolling writer without reading historical files.

  • MAP shared-shredding uses max-columns for the first file, then adjusts K from bounded statistics reported by completed files.
  • Adds an opt-in adaptive Variant inference mode. The first file uses the existing inference behavior, while subsequent files combine bounded prior evidence with a smaller current-file sample.
  • Uses admission and retention thresholds to avoid frequent Variant physical schema changes.
  • Keeps the existing per-file Variant inference as the default behavior.
  • Scopes adaptive state to one rolling writer and advances it only after a file is successfully completed.
  • Supports independent inference for multiple and nested Variant fields.
  • Validates adaptive inference sample sizes and thresholds.

Tests

  • Added unit tests for MAP column-count adaptation and Variant admission/retention behavior.
  • Added ORC and Parquet E2E tests for MAP shared-shredding across rolling files.
  • Added Parquet E2E tests for Variant rolling inference, multiple Variant fields, and nested Variant fields.
  • Added coverage for PK rolling-writer state isolation and inference option validation.
  • Relevant paimon-common and paimon-core tests pass, together with Spotless and diff checks.

@JingsongLi

Copy link
Copy Markdown
Contributor

Adaptive Cross-File Inference for Variant Shredding

1. Summary

Paimon currently infers a Variant shredding schema independently for every data file. Each file
buffers up to variant.shredding.maxInferBufferRow rows, infers its physical schema, flushes the
buffer, and then continues writing. This approach is robust to schema changes, but it repeats the
same inference work for adjacent files and may produce unnecessary physical-schema churn.

This document proposes an adaptive inference mode with three properties:

  1. Previous-file prior: retain bounded inference evidence and the selected physical schema from
    the preceding file in the same rolling writer.
  2. Small-sample correction: inspect a smaller prefix of the next file before finalizing its
    physical schema.
  3. Hysteresis: use different thresholds for retaining an existing shredded path and admitting a
    new path, reducing oscillation near the cardinality threshold.

The prior is an optimization hint, not a correctness dependency. Every file remains self-describing,
and values that do not fit a selected typed layout continue to be stored in the Variant value
field according to the Parquet Variant shredding model.

The state is intentionally scoped to one RollingFileWriter. It is not shared across tasks,
processes, partitions, buckets, or writer restarts, and it is not persisted in table metadata.

2. Motivation

Independent per-file inference has three costs.

2.1 Repeated buffering

For a stable input stream, adjacent files often infer almost identical schemas. Buffering the same
large prefix for every file adds memory pressure and delays creation of the underlying format
writer. It can also cause a file to exceed its target size because rolling cannot be evaluated until
the inferred plan has been finalized.

2.2 Physical-schema churn

A field whose observed cardinality is close to
variant.shredding.minFieldCardinalityRatio may be included in one file and excluded from the next.
The data remains readable, but frequent layout changes reduce predictability and can make projected
reads and data skipping less effective.

2.3 Weak inference for small files

Small files may close before enough rows have been observed. A bounded prior supplies evidence from
the immediately preceding stream segment without requiring table-wide coordination.

3. Goals

  • Reduce inference buffering after the first file in a rolling writer.
  • Preserve the ability to react to new fields and type changes at file boundaries.
  • Reduce repeated admission and eviction of fields near the selection threshold.
  • Preserve file-level correctness and independent readability.
  • Keep memory usage bounded by configured schema width and a fixed effective prior size.
  • Make the result deterministic for the same prior and current sample.
  • Avoid cross-writer synchronization and persistent table metadata in the first version.
  • Preserve the existing configured-schema behavior.

4. Non-Goals

  • Producing one globally stable shredding schema for an entire table.
  • Learning from files written by other tasks or writers.
  • Persisting inference state across checkpoints, retries, or process restarts.
  • Selecting fields from query history.
  • Changing the Variant encoding or Parquet shredding specification.
  • Changing the current strict type-merging policy in the first implementation.
  • Reconstructing prior statistics by reading the previous file footer.

Query-aware pinned paths and table-level schema pinning are complementary future features.

5. Current Behavior

InferShreddingWritePlanWriter buffers the prefix of each file and invokes
ShreddingWritePlanFactory.createWritePlan once the row limit is reached or the writer closes.
VariantShreddingWritePlanFactory creates a new InferVariantShreddingSchema and derives a
physical row type only from those buffered rows.

RowDataRollingFileWriter creates a new file-writer context for every rolled file. That operation
also creates a new format writer factory and therefore a new VariantShreddingWritePlanFactory.
No inference information survives the file boundary.

The current selection limits are:

  • cold sample limit: 4,096 rows by default;
  • minimum field cardinality ratio: 0.10;
  • maximum inferred schema width: 300;
  • maximum inference depth: 50.

6. Design Overview

Adaptive inference introduces one in-memory VariantShreddingInferenceSession for each sequential
rolling writer.

File 1 prefix (cold sample)
        |
        v
  analyze + select plan
        |
        +----> bounded evidence + selected plan
                              |
File 2 prefix (warm sample)   |
        |                     |
        +---------------------+
        |
        v
 combine evidence + apply hysteresis + correct types
        |
        +----> File 2 plan and updated bounded prior

The first file uses the existing cold-sample limit. Later files use a smaller warm-sample limit.
For every warm file, the session:

  1. analyzes the current prefix;
  2. combines current field-presence evidence with bounded prior evidence;
  3. applies admission and retention thresholds;
  4. validates the previous typed layout against current observations;
  5. applies the existing depth and width limits;
  6. creates the per-file write plan; and
  7. replaces the session state with a bounded posterior for the next file.

7. State Scope and Lifecycle

7.1 Scope

The session belongs to a single RollingFileWriter instance. In normal Paimon write paths, that
also limits it to one logical stream such as a partition/bucket/level writer.

The session must not be stored in:

  • the shared FileFormat, which is required to be thread-safe;
  • static state;
  • a table-wide cache; or
  • a cache keyed only by table or column name.

There is no globally meaningful "previous file" when multiple tasks write concurrently.

7.2 Lifecycle

  • The session is created with the rolling writer.
  • The first inferred plan initializes the session.
  • A subsequent file reads and updates the session only from the writer thread.
  • The rolling writer opens the next file only after the current file closes successfully.
  • If writing or closing fails, the rolling writer aborts and discards the entire session.
  • A retry creates a new session and starts with cold inference.
  • A configured shredding schema takes precedence and bypasses adaptive inference.

No extra commit protocol is required because the state is ephemeral and cannot be observed by
readers.

8. Inference Evidence

The implementation should separate analysis from schema materialization. The current
InferVariantShreddingSchema combines both operations, which makes it difficult to reuse bounded
evidence without carrying a partially materialized RowType.

A conceptual model is:

final class VariantShreddingInferenceState {
    private final Map<List<Integer>, VariantColumnEvidence> columns;
    private final RowType previousPhysicalRowType;
    private final long completedFileCount;
}

final class VariantColumnEvidence {
    private final double effectiveRootValueCount;
    private final Map<VariantFieldPath, FieldEvidence> fields;
}

final class FieldEvidence {
    private final double effectivePresenceCount;
    private final DataType previousSelectedType;
    private final boolean previouslySelected;
}

The actual implementation may use compact tree nodes instead of path-keyed maps. Field paths and
types must use deterministic equality and ordering.

For the first version, the reusable evidence only needs:

  • the effective number of non-null root Variant values;
  • the effective presence count of each object field;
  • whether the path was selected in the preceding plan; and
  • the type selected in the preceding plan.

The current sample still supplies exact type observations. Persisting a complete type histogram or
full-file fallback statistics can be added later.

8.1 Cardinality denominator

To preserve current behavior, a nested field's cardinality remains relative to the number of
non-null values of its root Variant column, rather than only to the number of occurrences of its
immediate parent.

Each Variant column maintains its own denominator. If the current prefix contains no non-null value
for a Variant column, that column's prior is retained without decay because the sample provides no
evidence about its shape.

9. Combining Prior and Current Evidence

Raw counts from the cold sample must not dominate all future files. The prior is therefore converted
to a bounded effective sample.

Let:

  • P be the previous bounded evidence;
  • C be evidence from the current warm prefix; and
  • W be the warm-sample row limit.

For each Variant column:

  1. scale P so that its effective root count is at most W;
  2. add the current counts from C;
  3. compute field presence ratios from the combined counts; and
  4. scale the combined evidence back to at most W before storing it as the next prior.

Conceptually:

boundedPrior = scaleToAtMost(previousEvidence, W)
posterior    = boundedPrior + currentEvidence
nextPrior    = scaleToAtMost(posterior, W)

With a full warm sample and the default equal cap, the previous evidence and the current file have
approximately equal influence. Repeated absence naturally decays a field's rate: 100%, 50%, 25%,
12.5%, and so on. The prior can therefore stabilize a path without making it permanent.

The first file may analyze up to the existing 4,096-row limit, but its evidence is scaled down before
it becomes a prior.

10. Hysteresis and Field Selection

Let:

  • A be the existing admission threshold,
    variant.shredding.minFieldCardinalityRatio; and
  • R be the retention threshold, with 0 <= R < A.

The proposed default is:

A = 0.10
R = 0.05

A field is eligible when:

previously selected: combinedPresenceRatio >= R
not previously selected: combinedPresenceRatio >= A

This gap prevents a field near 10% from alternating between shredded and unshredded layouts on
every file.

10.1 Width-budget ordering

Hysteresis must not allow stale fields to permanently occupy a full width budget. Eligible fields
are therefore ranked by:

  1. combined presence ratio, descending;
  2. previously selected status, selected first when ratios are equal; and
  3. complete field path, lexicographically ascending.

Parent nodes required to reach a selected child are included before the child. Existing width and
depth accounting remains authoritative. If two schemas consume the same budget, the deterministic
ordering produces the same result for the same evidence.

11. Type Selection and Correction

Field-membership hysteresis must not make type changes unnecessarily sticky.

For a path selected in the previous plan:

  1. If the current sample has no value for the path, retain the previous selected type.
  2. If all current observations are compatible with the previous type, retain it and apply existing
    decimal/integer widening rules when required.
  3. If any current observation is incompatible, discard the previous type hint and infer the type
    from the current sample using the existing strict merge rules.

For a newly admitted path, infer its type from the current sample using the existing rules.

This policy deliberately applies hysteresis to field membership, but immediate correction to
incompatible types. It avoids carrying a stale typed layout across several files while keeping
the first implementation behaviorally close to current Paimon.

Values after the sampled prefix may still fail to fit the selected type. They remain lossless in the
Variant value field. This is a correctness fallback, but a high fallback rate can reduce
compression, projection, and data-skipping benefits.

A follow-up design may select the most common compatible type with a configurable fit threshold,
similar to other Variant implementations. That policy change should be evaluated separately.

12. Buffer Limits

The adaptive mode uses:

  • the existing variant.shredding.maxInferBufferRow for the first file; and
  • a smaller warm-file row limit for subsequent files.

Inference should also stop at a byte limit. A row-only limit is insufficient when a small number of
Variant values contain large strings, objects, arrays, or binary values.

The writer should finalize its plan when either limit is reached:

bufferedRows >= rowLimit || estimatedBufferedBytes >= byteLimit

The byte estimator does not need to model Java object overhead exactly. It should conservatively
include Variant binary/value bytes and fixed row overhead, be deterministic, and avoid traversing
the same Variant tree twice.

13. Writer Integration

The smallest integration is to create the format writer factory once per
RowDataRollingFileWriter, then reuse that factory for the sequential files produced by the rolling
writer.

Today, RowDataRollingFileWriter calls:

fileFormat.createWriterFactory(writeSchema)

inside the per-file supplier. In adaptive mode, it should happen once when the supplier is built.
For Parquet, the resulting ShreddingWritePlanWriterFactory contains one
VariantShreddingWritePlanFactory, which in turn owns the adaptive session.

Per-file components must remain per-file:

  • output streams;
  • RowDataFileWriter;
  • statistics producers and collectors;
  • file-index writers;
  • sidecar writers; and
  • writer metadata.

Only the reusable FormatWriterFactory and its bounded inference session cross file boundaries.

ShreddingWritePlanFactory.inferBufferRowCount() may return the cold or warm limit based on whether
the session has a prior. createWritePlan(sampleRows) analyzes the current sample, materializes the
plan, and advances the session.

The factory is used sequentially by one rolling writer. It is not advertised as thread-safe and
must not be shared between rolling writers.

14. Configuration

The following options are proposed:

Option Default Description
variant.shredding.inferenceMode PER_FILE PER_FILE preserves current behavior; ADAPTIVE enables rolling-writer-scoped prior reuse.
variant.shredding.adaptive.maxInferBufferRow 256 Maximum prefix rows for files after the first file in the session.
variant.shredding.adaptive.retentionRatio 0.05 Minimum combined presence ratio for a path selected in the previous plan. Must not exceed variant.shredding.minFieldCardinalityRatio.
variant.shredding.maxInferBufferBytes 64 MB Maximum estimated bytes retained for inference in either mode.

The existing options continue to define:

  • cold-file row limit;
  • admission threshold;
  • maximum schema width; and
  • maximum schema depth.

The initial release should keep PER_FILE as the default. After production evidence demonstrates
stable write and read improvements, changing the default can be considered independently.

15. Correctness and Compatibility

15.1 Read correctness

Each file persists its own physical Variant shredding schema. Readers already support different
physical layouts across files. The adaptive session affects only how a file chooses that layout.

Values incompatible with a typed field are encoded in value, so prior reuse does not discard or
coerce user data.

15.2 Data skipping

A stale type may increase non-null value occurrences. Typed-column statistics alone cannot
describe those fallback values. The adaptive design therefore uses the current sample to correct
incompatible prior types immediately.

Future full-file fallback metrics should be used to detect paths whose typed layout is not
beneficial even when the prefix appeared compatible.

15.3 Backward compatibility

  • No file-format or table-metadata change is required.
  • Existing files require no rewrite.
  • Configured shredding schemas are unchanged.
  • PER_FILE mode preserves the current inference lifecycle.
  • Adaptive and non-adaptive files may coexist in the same table.

16. Failure, Retry, and Concurrency Semantics

  • Session mutations occur only in the rolling writer's thread.
  • The next file cannot consume the updated prior until the current file has closed.
  • If the current file fails, the writer aborts and no subsequent file observes its prior.
  • A task retry starts cold; its physical layouts may differ from the failed attempt, which is valid
    because layouts are file-local.
  • Concurrent writers never share inference sessions.
  • Schema evolution creates a new rolling writer/session. If an implementation can reuse a writer
    across logical schema changes, it must reset the session when the logical-row-type fingerprint
    changes.

17. Observability

Adaptive mode should expose metrics or debug counters for:

  • cold versus warm inference count;
  • sampled rows and estimated bytes;
  • inference time;
  • number of admitted, retained, evicted, and type-corrected paths;
  • inferred physical-schema width and depth;
  • physical-schema fingerprint changes between adjacent files; and
  • Variant values or bytes falling back from typed_value to value, if the writer can report them.

The schema fingerprint should be deterministic and suitable for measuring churn, but it does not
need to be persisted in table metadata.

18. Test Plan

18.1 Unit tests

Add focused tests for:

  • first-file cold inference and later-file warm limits;
  • stable adjacent files retaining the same plan;
  • a new hot field being admitted;
  • an old field decaying below the retention threshold;
  • a field oscillating between the admission and retention thresholds;
  • a full width budget replacing a stale low-frequency field with a hotter new field;
  • immediate correction of an incompatible previous type;
  • numeric widening compatible with the previous type;
  • no current non-null Variant values retaining the prior without decay;
  • independent state for multiple Variant columns;
  • nested fields, arrays, nulls, and depth limits;
  • deterministic ordering for equal-frequency paths; and
  • state reset after creating a new session.

18.2 Writer integration tests

Use a small target file size to force one rolling writer to create multiple Parquet files, then
verify:

  • the first and warm-file sample limits;
  • the physical schema stored in each file;
  • logical read equality across mixed physical layouts;
  • lossless fallback for values that appear after the sampled prefix;
  • behavior when a file fails or is aborted; and
  • no state leakage between two rolling writers.

18.3 Benchmarks

Compare:

  1. current per-file inference with a 4,096-row sample;
  2. direct previous-plan reuse without correction;
  3. adaptive prior plus a 256-row correction sample; and
  4. configured/pinned schema.

Use the following workloads:

  • stationary schemas;
  • sparse Zipf-distributed keys;
  • fields near the cardinality threshold;
  • abrupt schema shifts at file boundaries;
  • alternating types;
  • large Variant values;
  • partition- or sensor-sorted streams; and
  • compaction output mixing several input schemas.

Measure:

  • write throughput and time to initialize each physical writer;
  • peak inference-buffer bytes;
  • schema-fingerprint churn;
  • shredded-path coverage and fallback bytes;
  • file size and Parquet leaf-column count;
  • projected read bytes and latency; and
  • predicate/data-skipping effectiveness.

19. Rollout Plan

  1. Refactor Variant inference into analysis, evidence combination, and schema materialization
    without changing current behavior.
  2. Add the byte cap to per-file inference.
  3. Add the rolling-writer-scoped session behind ADAPTIVE mode.
  4. Add unit, integration, and benchmark coverage.
  5. Run shadow metrics that compare the adaptive plan with the current per-file plan without changing
    the written schema.
  6. Enable adaptive mode in selected workloads and evaluate schema churn, fallback rate, memory, and
    read performance.
  7. Consider changing the default only after results are stable across stationary and schema-shift
    workloads.

20. Alternatives Considered

20.1 Reuse the previous physical schema without sampling

This minimizes buffering but cannot detect new hot paths or boundary-aligned schema changes. It may
increase fallback values and weaken data skipping. It is not recommended.

20.2 Use a smaller independent sample for every file

This reduces buffering but loses evidence for stable sparse fields and increases threshold noise.
It does not address physical-schema churn.

20.3 Read the previous file footer

The footer exposes the selected physical schema but not the observations that caused the selection,
including frequency and rejected alternatives. The writer already has richer information in memory,
so footer I/O adds latency without improving the prior.

20.4 Persist one table-wide learned schema

This can provide stronger stability but requires coordination, evolution rules, ownership, and
metadata compatibility. It is better treated as explicit schema pinning rather than an implicit
extension of per-file inference.

20.5 Aggregate compaction input schemas

This is useful for compaction because there is no single representative previous input file.
Compaction-aware aggregation is compatible with this design but should be a separate extension.

21. Open Questions

  • Should the warm-sample default be 128, 256, or derived from the target file size?
  • Is a 0.05 absolute retention threshold preferable to a fixed fraction of the admission threshold?
  • Can the Variant writer report full-file fallback counts cheaply enough to feed the next prior?
  • Should compaction initialize its session from aggregated input-file metadata?
  • Which paths, if any, should users be able to pin ahead of inferred paths?
  • After the adaptive lifecycle is proven, should type selection move from strict consistency to a
    minimum fit ratio with minority values stored in value?

22. References

@JingsongLi

Copy link
Copy Markdown
Contributor

I used AI to create a design; I hope you find it useful.

@lxy-9602 lxy-9602 changed the title [core] Restore shredding write plans from historical files (variant & map shredding) [core] Support adaptive MAP and Variant shredding across rolling files Jul 24, 2026
@lxy-9602
lxy-9602 force-pushed the map-variant-restore branch from 906a92a to 189a50b Compare July 24, 2026 10:41
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants