This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 999df993ce [core] Support adaptive MAP and Variant shredding across 
rolling files (#8827)
999df993ce is described below

commit 999df993ceb21822ece95ce08321291882195d69
Author: lxy <[email protected]>
AuthorDate: Sun Jul 26 15:56:08 2026 +0800

    [core] Support adaptive MAP and Variant shredding across rolling files 
(#8827)
---
 docs/generated/core_configuration.html             |  18 +
 .../main/java/org/apache/paimon/CoreOptions.java   |  55 +++
 .../shredding/MapSharedShreddingWritePlan.java     |   8 +
 .../MapSharedShreddingWritePlanFactory.java        |  46 +--
 .../data/variant/InferVariantShreddingSchema.java  | 387 +++++++++++++++++++--
 .../variant/VariantShreddingInferenceSession.java  |  91 +++++
 .../shredding/InferShreddingWritePlanWriter.java   |   2 +-
 .../format/shredding/ShreddingFormatWriter.java    |   4 +
 .../shredding/ShreddingWritePlanFactory.java       |   8 +
 .../shredding/ShreddingWritePlanWriterFactory.java |   3 +
 .../variant/VariantShreddingWritePlanFactory.java  |  53 ++-
 .../shredding/MapSharedShreddingWritePlanTest.java |  31 +-
 .../variant/InferVariantShreddingSchemaTest.java   |  97 ++++++
 .../shredding/ShreddingFormatWriterTest.java       |  24 +-
 .../paimon/io/KeyValueFileWriterFactory.java       |  47 ++-
 .../apache/paimon/io/RowDataRollingFileWriter.java |  58 +--
 .../org/apache/paimon/schema/SchemaValidation.java |  39 +++
 .../apache/paimon/append/AppendOnlyWriterTest.java |  23 +-
 .../apache/paimon/schema/SchemaValidationTest.java |  77 ++++
 .../paimon/table/MapSharedShreddingTableTest.java  |  34 +-
 .../writer/InferVariantShreddingWriteTest.java     | 303 +++++++++++++++-
 21 files changed, 1275 insertions(+), 133 deletions(-)

diff --git a/docs/generated/core_configuration.html 
b/docs/generated/core_configuration.html
index 630b6eb200..b1da97ffea 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -1735,6 +1735,24 @@ If the data size allocated for the sorting task is 
uneven,which may lead to perf
             <td>Boolean</td>
             <td>Whether to automatically infer the shredding schema when 
writing Variant columns.</td>
         </tr>
+        <tr>
+            <td><h5>variant.shredding.adaptive.maxInferBufferRow</h5></td>
+            <td style="word-wrap: break-word;">256</td>
+            <td>Integer</td>
+            <td>Maximum number of prefix rows sampled after the first file in 
an adaptive Variant shredding inference session.</td>
+        </tr>
+        <tr>
+            <td><h5>variant.shredding.adaptive.retentionRatio</h5></td>
+            <td style="word-wrap: break-word;">0.05</td>
+            <td>Double</td>
+            <td>Minimum combined presence ratio for retaining a Variant path 
selected in the previous file. This must not exceed 
'variant.shredding.minFieldCardinalityRatio'.</td>
+        </tr>
+        <tr>
+            <td><h5>variant.shredding.inferenceMode</h5></td>
+            <td style="word-wrap: break-word;">per-file</td>
+            <td><p>Enum</p></td>
+            <td>The Variant shredding inference mode. PER_FILE infers each 
file independently. ADAPTIVE reuses bounded evidence within one rolling writer 
and samples a smaller prefix after the first file.<br /><br />Possible 
values:<ul><li>"per-file": Infer every file independently from its own prefix 
rows.</li><li>"adaptive": Reuse bounded inference evidence within one rolling 
writer and correct it with a smaller prefix from each subsequent 
file.</li></ul></td>
+        </tr>
         <tr>
             <td><h5>variant.shredding.maxInferBufferRow</h5></td>
             <td style="word-wrap: break-word;">4096</td>
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java 
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index 6a06832770..2def6e9ffc 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -453,6 +453,17 @@ public class CoreOptions implements Serializable {
                     .withDescription(
                             "Whether to automatically infer the shredding 
schema when writing Variant columns.");
 
+    public static final ConfigOption<VariantShreddingInferenceMode>
+            VARIANT_SHREDDING_INFERENCE_MODE =
+                    key("variant.shredding.inferenceMode")
+                            .enumType(VariantShreddingInferenceMode.class)
+                            
.defaultValue(VariantShreddingInferenceMode.PER_FILE)
+                            .withDescription(
+                                    "The Variant shredding inference mode. 
PER_FILE infers each "
+                                            + "file independently. ADAPTIVE 
reuses bounded evidence "
+                                            + "within one rolling writer and 
samples a smaller "
+                                            + "prefix after the first file.");
+
     public static final ConfigOption<Integer> 
VARIANT_SHREDDING_MAX_SCHEMA_WIDTH =
             key("variant.shredding.maxSchemaWidth")
                     .intType()
@@ -481,6 +492,23 @@ public class CoreOptions implements Serializable {
                     .defaultValue(4096)
                     .withDescription("Maximum number of rows to buffer for 
schema inference.");
 
+    public static final ConfigOption<Integer> 
VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW =
+            key("variant.shredding.adaptive.maxInferBufferRow")
+                    .intType()
+                    .defaultValue(256)
+                    .withDescription(
+                            "Maximum number of prefix rows sampled after the 
first file in "
+                                    + "an adaptive Variant shredding inference 
session.");
+
+    public static final ConfigOption<Double> 
VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO =
+            key("variant.shredding.adaptive.retentionRatio")
+                    .doubleType()
+                    .defaultValue(0.05)
+                    .withDescription(
+                            "Minimum combined presence ratio for retaining a 
Variant path selected "
+                                    + "in the previous file. This must not 
exceed "
+                                    + 
"'variant.shredding.minFieldCardinalityRatio'.");
+
     public static final ConfigOption<String> MANIFEST_FORMAT =
             key("manifest.format")
                     .stringType()
@@ -5483,6 +5511,33 @@ public class CoreOptions implements Serializable {
         }
     }
 
+    /** Inference mode for Variant shredding schemas. */
+    public enum VariantShreddingInferenceMode implements DescribedEnum {
+        PER_FILE("per-file", "Infer every file independently from its own 
prefix rows."),
+        ADAPTIVE(
+                "adaptive",
+                "Reuse bounded inference evidence within one rolling writer 
and correct it with "
+                        + "a smaller prefix from each subsequent file.");
+
+        private final String value;
+        private final String description;
+
+        VariantShreddingInferenceMode(String value, String description) {
+            this.value = value;
+            this.description = description;
+        }
+
+        @Override
+        public String toString() {
+            return value;
+        }
+
+        @Override
+        public InlineElement getDescription() {
+            return text(description);
+        }
+    }
+
     /**
      * Action to take when an UPDATE (e.g. via MERGE INTO) modifies columns 
that are covered by a
      * global index.
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlan.java
 
b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlan.java
index 9555a2f45e..8b620d92df 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlan.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlan.java
@@ -78,4 +78,12 @@ public class MapSharedShreddingWritePlan implements 
ShreddingWritePlan {
         }
         return Collections.unmodifiableMap(metadata);
     }
+
+    Map<String, Integer> fileMaxRowWidths() {
+        Map<String, Integer> result = new LinkedHashMap<>();
+        for (String fieldName : converter.shreddingFieldNames()) {
+            result.put(fieldName, 
converter.buildFieldMeta(fieldName).maxRowWidth());
+        }
+        return result;
+    }
 }
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanFactory.java
 
b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanFactory.java
index 7818b24baa..427e6cb9ef 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanFactory.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanFactory.java
@@ -20,7 +20,6 @@ package org.apache.paimon.data.shredding;
 
 import org.apache.paimon.CoreOptions;
 import org.apache.paimon.CoreOptions.MapSharedShreddingColumnPlacementPolicy;
-import org.apache.paimon.data.InternalMap;
 import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.format.shredding.ShreddingWritePlanFactory;
 import org.apache.paimon.options.Options;
@@ -35,12 +34,10 @@ import static 
org.apache.paimon.utils.Preconditions.checkArgument;
 /** Creates per-file shared-shredding MAP write plans. */
 public class MapSharedShreddingWritePlanFactory implements 
ShreddingWritePlanFactory {
 
-    private static final int INFER_BUFFER_ROW_COUNT = 1;
-
     private final RowType logicalRowType;
     private final Map<String, Integer> fieldToMaxColumns;
     private final Map<String, MapSharedShreddingColumnPlacementPolicy> 
fieldToColumnPlacementPolicy;
-    private final Map<String, Integer> fieldToPosition;
+    private final MapSharedShreddingContext context;
 
     public MapSharedShreddingWritePlanFactory(RowType logicalRowType, Options 
options) {
         this.logicalRowType = logicalRowType;
@@ -50,12 +47,11 @@ public class MapSharedShreddingWritePlanFactory implements 
ShreddingWritePlanFac
         this.fieldToMaxColumns =
                 
MapSharedShreddingUtils.buildColumnToNumColumns(shreddingFields, coreOptions);
         this.fieldToColumnPlacementPolicy = new LinkedHashMap<>();
-        this.fieldToPosition = new LinkedHashMap<>();
         for (String field : shreddingFields) {
             fieldToColumnPlacementPolicy.put(
                     field, 
coreOptions.mapSharedShreddingColumnPlacementPolicy(field));
-            fieldToPosition.put(field, logicalRowType.getFieldIndex(field));
         }
+        this.context = new MapSharedShreddingContext(fieldToMaxColumns);
     }
 
     @Override
@@ -70,39 +66,29 @@ public class MapSharedShreddingWritePlanFactory implements 
ShreddingWritePlanFac
 
     @Override
     public boolean shouldInferWritePlan() {
-        return shouldCreateWritePlan();
+        return false;
     }
 
     @Override
     public int inferBufferRowCount() {
-        return INFER_BUFFER_ROW_COUNT;
+        return 0;
     }
 
     @Override
     public ShreddingWritePlan createWritePlan(List<InternalRow> sampleRows) {
         checkArgument(shouldCreateWritePlan(), "MAP shared-shredding write 
plan is not active.");
-
-        Map<String, Integer> fieldToNumColumns = new LinkedHashMap<>();
-        for (Map.Entry<String, Integer> entry : fieldToMaxColumns.entrySet()) {
-            int maxColumns = entry.getValue();
-            int numColumns = maxColumns;
-            if (!sampleRows.isEmpty()) {
-                int fieldPosition = fieldToPosition.get(entry.getKey());
-                int maxRowWidth = 0;
-                int sampleCount = Math.min(sampleRows.size(), 
INFER_BUFFER_ROW_COUNT);
-                for (int i = 0; i < sampleCount; i++) {
-                    InternalRow row = sampleRows.get(i);
-                    InternalMap map =
-                            row.isNullAt(fieldPosition) ? null : 
row.getMap(fieldPosition);
-                    maxRowWidth = Math.max(maxRowWidth, map == null ? 0 : 
map.size());
-                }
-                numColumns = Math.max(1, Math.min(maxRowWidth, maxColumns));
-            }
-            fieldToNumColumns.put(entry.getKey(), numColumns);
-        }
-
-        // TODO: Infer the column count from recent file metadata instead of 
current-file samples.
         return new MapSharedShreddingWritePlan(
-                logicalRowType, fieldToNumColumns, 
fieldToColumnPlacementPolicy);
+                logicalRowType, context.computeNextK(), 
fieldToColumnPlacementPolicy);
+    }
+
+    @Override
+    public void onFileCompleted(ShreddingWritePlan writePlan) {
+        checkArgument(
+                writePlan instanceof MapSharedShreddingWritePlan,
+                "Unexpected MAP shared-shredding write plan: %s",
+                writePlan.getClass().getName());
+        ((MapSharedShreddingWritePlan) writePlan)
+                .fileMaxRowWidths()
+                .forEach(context::reportFileStats);
     }
 }
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/data/variant/InferVariantShreddingSchema.java
 
b/paimon-common/src/main/java/org/apache/paimon/data/variant/InferVariantShreddingSchema.java
index 94e268aed5..f18aaba0ff 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/data/variant/InferVariantShreddingSchema.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/data/variant/InferVariantShreddingSchema.java
@@ -29,6 +29,7 @@ import org.apache.paimon.types.VariantType;
 
 import java.math.BigDecimal;
 import java.util.ArrayList;
+import java.util.Comparator;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -60,34 +61,164 @@ public class InferVariantShreddingSchema {
 
     /** Infer schema from a list of rows. */
     public RowType inferSchema(List<InternalRow> rows) {
+        return inferInitial(rows, Integer.MAX_VALUE).physicalRowType();
+    }
+
+    AdaptiveInferenceResult inferInitial(List<InternalRow> rows, int 
effectiveSampleSize) {
+        InferenceEvidence evidence = analyze(rows);
         MaxFields maxFields = new MaxFields(maxSchemaWidth);
         Map<List<Integer>, RowType> inferredSchemas = new HashMap<>();
+        Map<List<Integer>, DataType> selectedSchemas = new HashMap<>();
 
         for (List<Integer> path : pathsToVariant) {
-            int numNonNullValues = 0;
-            DataType simpleSchema = null;
-
-            for (InternalRow row : rows) {
-                Variant variant = getValueAtPath(schema, row, path);
-                if (variant != null) {
-                    numNonNullValues++;
-                    GenericVariant v = (GenericVariant) variant;
-                    DataType schemaOfRow = schemaOf(v, maxSchemaDepth);
-                    simpleSchema = mergeSchema(simpleSchema, schemaOfRow);
-                }
-            }
+            ColumnEvidence column = evidence.columns.get(path);
+            double numNonNullValues = column == null ? 0 : 
column.rootValueCount;
+            DataType simpleSchema = column == null ? null : 
column.observedSchema;
 
             // Don't infer a schema for fields that appear in less than 
minFieldCardinalityRatio
             int minCardinality = (int) Math.ceil(numNonNullValues * 
minFieldCardinalityRatio);
 
             DataType finalizedSchema =
                     finalizeSimpleSchema(simpleSchema, minCardinality, 
maxFields);
+            selectedSchemas.put(path, finalizedSchema);
             RowType shreddingSchema = 
PaimonShreddingUtils.variantShreddingSchema(finalizedSchema);
             inferredSchemas.put(path, shreddingSchema);
         }
 
         // Insert each inferred schema into the full schema
-        return updateSchema(schema, inferredSchemas, new ArrayList<>());
+        return new AdaptiveInferenceResult(
+                updateSchema(schema, inferredSchemas, new ArrayList<>()),
+                limitEvidence(evidence, effectiveSampleSize),
+                selectedSchemas);
+    }
+
+    InferenceEvidence analyze(List<InternalRow> rows) {
+        Map<List<Integer>, ColumnEvidence> columns = new HashMap<>();
+        for (List<Integer> path : pathsToVariant) {
+            double rootValueCount = 0;
+            DataType observedSchema = null;
+            for (InternalRow row : rows) {
+                Variant variant = getValueAtPath(schema, row, path);
+                if (variant != null) {
+                    rootValueCount++;
+                    observedSchema =
+                            mergeSchema(
+                                    observedSchema,
+                                    schemaOf((GenericVariant) variant, 
maxSchemaDepth));
+                }
+            }
+            columns.put(path, new ColumnEvidence(rootValueCount, 
observedSchema));
+        }
+        return new InferenceEvidence(columns);
+    }
+
+    AdaptiveInferenceResult inferAdaptive(
+            InferenceEvidence previousEvidence,
+            Map<List<Integer>, DataType> previousSelectedSchemas,
+            List<InternalRow> rows,
+            int effectiveSampleSize,
+            double admissionRatio,
+            double retentionRatio) {
+        InferenceEvidence currentEvidence = analyze(rows);
+        InferenceEvidence combinedEvidence =
+                combineEvidence(previousEvidence, currentEvidence, 
effectiveSampleSize);
+        MaxFields maxFields = new MaxFields(maxSchemaWidth);
+        Map<List<Integer>, RowType> shreddingSchemas = new HashMap<>();
+        Map<List<Integer>, DataType> selectedSchemas = new HashMap<>();
+
+        for (List<Integer> path : pathsToVariant) {
+            ColumnEvidence combined = combinedEvidence.columns.get(path);
+            ColumnEvidence current = currentEvidence.columns.get(path);
+            DataType previousSelected = previousSelectedSchemas.get(path);
+            DataType selected =
+                    finalizeAdaptiveSchema(
+                            combined == null ? null : combined.observedSchema,
+                            current == null ? null : current.observedSchema,
+                            previousSelected,
+                            combined == null ? 0 : combined.rootValueCount,
+                            admissionRatio,
+                            retentionRatio,
+                            maxFields);
+            selectedSchemas.put(path, selected);
+            shreddingSchemas.put(path, 
PaimonShreddingUtils.variantShreddingSchema(selected));
+        }
+
+        return new AdaptiveInferenceResult(
+                updateSchema(schema, shreddingSchemas, new ArrayList<>()),
+                combinedEvidence,
+                selectedSchemas);
+    }
+
+    private InferenceEvidence combineEvidence(
+            InferenceEvidence previous, InferenceEvidence current, int 
effectiveSampleSize) {
+        Map<List<Integer>, ColumnEvidence> columns = new HashMap<>();
+        for (List<Integer> path : pathsToVariant) {
+            ColumnEvidence previousColumn = previous == null ? null : 
previous.columns.get(path);
+            ColumnEvidence currentColumn = current.columns.get(path);
+            if (currentColumn == null || currentColumn.rootValueCount == 0) {
+                if (previousColumn != null) {
+                    columns.put(path, previousColumn);
+                } else {
+                    columns.put(path, new ColumnEvidence(0, null));
+                }
+                continue;
+            }
+
+            ColumnEvidence boundedPrevious = scaleToAtMost(previousColumn, 
effectiveSampleSize);
+            double rootValueCount =
+                    currentColumn.rootValueCount
+                            + (boundedPrevious == null ? 0 : 
boundedPrevious.rootValueCount);
+            DataType observedSchema =
+                    mergeSchema(
+                            boundedPrevious == null ? null : 
boundedPrevious.observedSchema,
+                            currentColumn.observedSchema);
+            columns.put(
+                    path,
+                    scaleToAtMost(
+                            new ColumnEvidence(rootValueCount, observedSchema),
+                            effectiveSampleSize));
+        }
+        return new InferenceEvidence(columns);
+    }
+
+    private InferenceEvidence limitEvidence(InferenceEvidence evidence, int 
effectiveSampleSize) {
+        Map<List<Integer>, ColumnEvidence> columns = new HashMap<>();
+        for (Map.Entry<List<Integer>, ColumnEvidence> entry : 
evidence.columns.entrySet()) {
+            columns.put(entry.getKey(), scaleToAtMost(entry.getValue(), 
effectiveSampleSize));
+        }
+        return new InferenceEvidence(columns);
+    }
+
+    private ColumnEvidence scaleToAtMost(ColumnEvidence evidence, int 
maxRootValueCount) {
+        if (evidence == null
+                || evidence.rootValueCount == 0
+                || evidence.rootValueCount <= maxRootValueCount) {
+            return evidence;
+        }
+        double scale = maxRootValueCount / evidence.rootValueCount;
+        return new ColumnEvidence(
+                maxRootValueCount, scaleFieldCounts(evidence.observedSchema, 
scale));
+    }
+
+    private DataType scaleFieldCounts(DataType dataType, double scale) {
+        if (dataType instanceof RowType) {
+            List<DataField> fields = new ArrayList<>();
+            for (DataField field : ((RowType) dataType).getFields()) {
+                fields.add(
+                        new DataField(
+                                field.id(),
+                                field.name(),
+                                scaleFieldCounts(field.type(), scale),
+                                String.valueOf(getFieldCount(field) * scale)));
+            }
+            return new RowType(dataType.isNullable(), fields);
+        }
+        if (dataType instanceof ArrayType) {
+            ArrayType arrayType = (ArrayType) dataType;
+            return new ArrayType(
+                    arrayType.isNullable(), 
scaleFieldCounts(arrayType.getElementType(), scale));
+        }
+        return dataType;
     }
 
     /**
@@ -255,7 +386,7 @@ public class InferVariantShreddingSchema {
         }
     }
 
-    private long getFieldCount(DataField field) {
+    private double getFieldCount(DataField field) {
         // Read count from description field
         String desc = field.description();
         if (desc == null || desc.isEmpty()) {
@@ -264,7 +395,7 @@ public class InferVariantShreddingSchema {
                             "Field '%s' is missing count in description. This 
should not happen during schema inference.",
                             field.name()));
         }
-        return Long.parseLong(desc);
+        return Double.parseDouble(desc);
     }
 
     /** Merge two decimals with possibly different scales. */
@@ -338,8 +469,8 @@ public class InferVariantShreddingSchema {
 
             if (comp == 0) {
                 DataType dataType = mergeSchema(field1.type(), field2.type());
-                long c1 = getFieldCount(field1);
-                long c2 = getFieldCount(field2);
+                double c1 = getFieldCount(field1);
+                double c2 = getFieldCount(field2);
                 // Store count in description
                 DataField newField =
                         new DataField(nextFieldId++, f1Name, dataType, 
String.valueOf(c1 + c2));
@@ -347,14 +478,14 @@ public class InferVariantShreddingSchema {
                 f1Idx++;
                 f2Idx++;
             } else if (comp < 0) {
-                long count = getFieldCount(field1);
+                double count = getFieldCount(field1);
                 DataField newField =
                         new DataField(
                                 nextFieldId++, field1.name(), field1.type(), 
String.valueOf(count));
                 newFields.add(newField);
                 f1Idx++;
             } else {
-                long count = getFieldCount(field2);
+                double count = getFieldCount(field2);
                 DataField newField =
                         new DataField(
                                 nextFieldId++, field2.name(), field2.type(), 
String.valueOf(count));
@@ -365,7 +496,7 @@ public class InferVariantShreddingSchema {
 
         while (f1Idx < fields1.size() && newFields.size() < maxRowFieldSize) {
             DataField field1 = fields1.get(f1Idx);
-            long count = getFieldCount(field1);
+            double count = getFieldCount(field1);
             DataField newField =
                     new DataField(
                             nextFieldId++, field1.name(), field1.type(), 
String.valueOf(count));
@@ -375,7 +506,7 @@ public class InferVariantShreddingSchema {
 
         while (f2Idx < fields2.size() && newFields.size() < maxRowFieldSize) {
             DataField field2 = fields2.get(f2Idx);
-            long count = getFieldCount(field2);
+            double count = getFieldCount(field2);
             DataField newField =
                     new DataField(
                             nextFieldId++, field2.name(), field2.type(), 
String.valueOf(count));
@@ -420,6 +551,172 @@ public class InferVariantShreddingSchema {
         return new RowType(newFields);
     }
 
+    private DataType finalizeAdaptiveSchema(
+            DataType combined,
+            DataType current,
+            DataType previousSelected,
+            double rootValueCount,
+            double admissionRatio,
+            double retentionRatio,
+            MaxFields maxFields) {
+        maxFields.remaining--;
+        if (maxFields.remaining <= 0) {
+            return DataTypes.VARIANT();
+        }
+
+        if (current != null
+                && previousSelected != null
+                && !compatibleTypeFamilies(previousSelected, current)) {
+            combined = current;
+            previousSelected = null;
+        }
+        if (combined == null || combined instanceof VariantType) {
+            if (current != null && !(current instanceof VariantType)) {
+                combined = current;
+            } else if (previousSelected != null) {
+                combined = previousSelected;
+            } else {
+                return DataTypes.VARIANT();
+            }
+        }
+
+        if (combined instanceof RowType) {
+            RowType combinedRow = (RowType) combined;
+            RowType currentRow = current instanceof RowType ? (RowType) 
current : null;
+            RowType previousRow =
+                    previousSelected instanceof RowType ? (RowType) 
previousSelected : null;
+            List<DataField> candidates = new ArrayList<>();
+            for (DataField field : combinedRow.getFields()) {
+                DataField previousField = findField(previousRow, field.name());
+                double threshold = previousField == null ? admissionRatio : 
retentionRatio;
+                double ratio = rootValueCount == 0 ? 0 : getFieldCount(field) 
/ rootValueCount;
+                if (ratio >= threshold) {
+                    candidates.add(field);
+                }
+            }
+            candidates.sort(
+                    Comparator.comparingDouble(
+                                    (DataField field) ->
+                                            rootValueCount == 0
+                                                    ? 0
+                                                    : getFieldCount(field) / 
rootValueCount)
+                            .reversed()
+                            .thenComparing(field -> findField(previousRow, 
field.name()) == null)
+                            .thenComparing(DataField::name));
+
+            List<DataField> selected = new ArrayList<>();
+            for (DataField field : candidates) {
+                if (maxFields.remaining <= 0) {
+                    break;
+                }
+                DataField currentField = findField(currentRow, field.name());
+                DataField previousField = findField(previousRow, field.name());
+                DataType selectedType =
+                        finalizeAdaptiveSchema(
+                                field.type(),
+                                currentField == null ? null : 
currentField.type(),
+                                previousField == null ? null : 
previousField.type(),
+                                rootValueCount,
+                                admissionRatio,
+                                retentionRatio,
+                                maxFields);
+                selected.add(new DataField(0, field.name(), selectedType));
+            }
+            selected.sort(Comparator.comparing(DataField::name));
+            List<DataField> fields = new ArrayList<>();
+            for (int i = 0; i < selected.size(); i++) {
+                DataField field = selected.get(i);
+                fields.add(new DataField(i, field.name(), field.type()));
+            }
+            return fields.isEmpty() ? DataTypes.VARIANT() : new 
RowType(fields);
+        }
+
+        if (combined instanceof ArrayType) {
+            ArrayType combinedArray = (ArrayType) combined;
+            DataType currentElement =
+                    current instanceof ArrayType ? ((ArrayType) 
current).getElementType() : null;
+            DataType previousElement =
+                    previousSelected instanceof ArrayType
+                            ? ((ArrayType) previousSelected).getElementType()
+                            : null;
+            return new ArrayType(
+                    finalizeAdaptiveSchema(
+                            combinedArray.getElementType(),
+                            currentElement,
+                            previousElement,
+                            rootValueCount,
+                            admissionRatio,
+                            retentionRatio,
+                            maxFields));
+        }
+
+        maxFields.remaining--;
+        return selectScalarType(combined, current, previousSelected);
+    }
+
+    private DataType selectScalarType(
+            DataType combined, DataType current, DataType previousSelected) {
+        if (current == null) {
+            return previousSelected == null ? widenScalarType(combined) : 
previousSelected;
+        }
+        if (previousSelected == null) {
+            return widenScalarType(current);
+        }
+
+        DataType merged = mergeScalarTypes(previousSelected, current);
+        return merged instanceof VariantType ? widenScalarType(current) : 
merged;
+    }
+
+    private boolean compatibleTypeFamilies(DataType previous, DataType 
current) {
+        if (previous instanceof RowType || current instanceof RowType) {
+            return previous instanceof RowType && current instanceof RowType;
+        }
+        if (previous instanceof ArrayType || current instanceof ArrayType) {
+            return previous instanceof ArrayType && current instanceof 
ArrayType;
+        }
+        return !(mergeScalarTypes(previous, current) instanceof VariantType);
+    }
+
+    private DataType mergeScalarTypes(DataType first, DataType second) {
+        if (first instanceof DecimalType && second instanceof DecimalType) {
+            return mergeDecimal((DecimalType) first, (DecimalType) second);
+        }
+        if (first instanceof DecimalType
+                && second.getTypeRoot() == 
org.apache.paimon.types.DataTypeRoot.BIGINT) {
+            return mergeDecimalWithLong((DecimalType) first);
+        }
+        if (first.getTypeRoot() == org.apache.paimon.types.DataTypeRoot.BIGINT
+                && second instanceof DecimalType) {
+            return mergeDecimalWithLong((DecimalType) second);
+        }
+        return first.equals(second) ? first : DataTypes.VARIANT();
+    }
+
+    private DataType widenScalarType(DataType dataType) {
+        if (dataType instanceof DecimalType) {
+            DecimalType decimalType = (DecimalType) dataType;
+            if (decimalType.getPrecision() <= 18 && decimalType.getScale() == 
0) {
+                return DataTypes.BIGINT();
+            }
+            return decimalType.getPrecision() <= 18
+                    ? DataTypes.DECIMAL(18, decimalType.getScale())
+                    : DataTypes.DECIMAL(DecimalType.MAX_PRECISION, 
decimalType.getScale());
+        }
+        return dataType;
+    }
+
+    private DataField findField(RowType rowType, String fieldName) {
+        if (rowType == null) {
+            return null;
+        }
+        for (DataField field : rowType.getFields()) {
+            if (field.name().equals(fieldName)) {
+                return field;
+            }
+        }
+        return null;
+    }
+
     /** Container for a mutable integer to track the total number of shredded 
fields. */
     private static class MaxFields {
         int remaining;
@@ -429,6 +726,54 @@ public class InferVariantShreddingSchema {
         }
     }
 
+    static final class InferenceEvidence {
+
+        private final Map<List<Integer>, ColumnEvidence> columns;
+
+        private InferenceEvidence(Map<List<Integer>, ColumnEvidence> columns) {
+            this.columns = columns;
+        }
+    }
+
+    private static final class ColumnEvidence {
+
+        private final double rootValueCount;
+        private final DataType observedSchema;
+
+        private ColumnEvidence(double rootValueCount, DataType observedSchema) 
{
+            this.rootValueCount = rootValueCount;
+            this.observedSchema = observedSchema;
+        }
+    }
+
+    static final class AdaptiveInferenceResult {
+
+        private final RowType physicalRowType;
+        private final InferenceEvidence evidence;
+        private final Map<List<Integer>, DataType> selectedSchemas;
+
+        private AdaptiveInferenceResult(
+                RowType physicalRowType,
+                InferenceEvidence evidence,
+                Map<List<Integer>, DataType> selectedSchemas) {
+            this.physicalRowType = physicalRowType;
+            this.evidence = evidence;
+            this.selectedSchemas = selectedSchemas;
+        }
+
+        RowType physicalRowType() {
+            return physicalRowType;
+        }
+
+        InferenceEvidence evidence() {
+            return evidence;
+        }
+
+        Map<List<Integer>, DataType> selectedSchemas() {
+            return selectedSchemas;
+        }
+    }
+
     /**
      * Given the schema of a Variant type, finalize the schema. Specifically: 
1) Widen integer types
      * to LongType 2) Replace empty rows with VariantType 3) Limit the total 
number of shredded
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/data/variant/VariantShreddingInferenceSession.java
 
b/paimon-common/src/main/java/org/apache/paimon/data/variant/VariantShreddingInferenceSession.java
new file mode 100644
index 0000000000..ba4292d378
--- /dev/null
+++ 
b/paimon-common/src/main/java/org/apache/paimon/data/variant/VariantShreddingInferenceSession.java
@@ -0,0 +1,91 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.data.variant;
+
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.Preconditions;
+
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Rolling-writer-scoped Variant inference state.
+ *
+ * <p>The state is intentionally not persisted or shared between rolling 
writers.
+ */
+public class VariantShreddingInferenceSession {
+
+    private final InferVariantShreddingSchema inferrer;
+    private final int effectiveSampleSize;
+    private final double admissionRatio;
+    private final double retentionRatio;
+
+    private InferVariantShreddingSchema.InferenceEvidence committedEvidence;
+    private Map<List<Integer>, DataType> committedSelectedSchemas = 
Collections.emptyMap();
+    private InferVariantShreddingSchema.AdaptiveInferenceResult pendingResult;
+
+    public VariantShreddingInferenceSession(
+            InferVariantShreddingSchema inferrer,
+            int effectiveSampleSize,
+            double admissionRatio,
+            double retentionRatio) {
+        Preconditions.checkArgument(
+                effectiveSampleSize > 0, "Effective sample size must be 
positive.");
+        Preconditions.checkArgument(
+                admissionRatio >= 0 && admissionRatio <= 1,
+                "Admission ratio must be between 0 and 1.");
+        Preconditions.checkArgument(
+                retentionRatio >= 0 && retentionRatio <= admissionRatio,
+                "Retention ratio must be between 0 and the admission ratio.");
+        this.inferrer = inferrer;
+        this.effectiveSampleSize = effectiveSampleSize;
+        this.admissionRatio = admissionRatio;
+        this.retentionRatio = retentionRatio;
+    }
+
+    public boolean hasPrior() {
+        return committedEvidence != null;
+    }
+
+    public RowType inferSchema(List<InternalRow> rows) {
+        if (committedEvidence == null) {
+            pendingResult = inferrer.inferInitial(rows, effectiveSampleSize);
+        } else {
+            pendingResult =
+                    inferrer.inferAdaptive(
+                            committedEvidence,
+                            committedSelectedSchemas,
+                            rows,
+                            effectiveSampleSize,
+                            admissionRatio,
+                            retentionRatio);
+        }
+        return pendingResult.physicalRowType();
+    }
+
+    public void commitPendingInference() {
+        Preconditions.checkState(pendingResult != null, "No pending Variant 
inference to commit.");
+        committedEvidence = pendingResult.evidence();
+        committedSelectedSchemas = pendingResult.selectedSchemas();
+        pendingResult = null;
+    }
+}
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/format/shredding/InferShreddingWritePlanWriter.java
 
b/paimon-common/src/main/java/org/apache/paimon/format/shredding/InferShreddingWritePlanWriter.java
index 5f497d2d92..22f102d814 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/format/shredding/InferShreddingWritePlanWriter.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/format/shredding/InferShreddingWritePlanWriter.java
@@ -126,7 +126,7 @@ public class InferShreddingWritePlanWriter implements 
BundleFormatWriter {
         ShreddingWritePlan writePlan = 
writePlanFactory.createWritePlan(collectAllRows());
         actualWriter =
                 ShreddingWritePlanWriterFactory.createWriterWithPlan(
-                        writerFactory, out, compression, writePlan);
+                        writerFactory, writePlanFactory, out, compression, 
writePlan);
         planFinalized = true;
 
         if (!bufferedBundles.isEmpty()) {
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingFormatWriter.java
 
b/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingFormatWriter.java
index 504852819f..97a9d4d2da 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingFormatWriter.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingFormatWriter.java
@@ -35,16 +35,19 @@ public class ShreddingFormatWriter implements 
BundleFormatWriter {
 
     private final FormatWriter delegate;
     private final SupportsShreddingWritePlan writerFactory;
+    private final ShreddingWritePlanFactory writePlanFactory;
     private final ShreddingWritePlan writePlan;
     private final String compression;
 
     public ShreddingFormatWriter(
             FormatWriter delegate,
             SupportsShreddingWritePlan writerFactory,
+            ShreddingWritePlanFactory writePlanFactory,
             ShreddingWritePlan writePlan,
             String compression) {
         this.delegate = delegate;
         this.writerFactory = writerFactory;
+        this.writePlanFactory = writePlanFactory;
         this.writePlan = writePlan;
         this.compression = compression;
     }
@@ -84,6 +87,7 @@ public class ShreddingFormatWriter implements 
BundleFormatWriter {
         } finally {
             delegate.close();
         }
+        writePlanFactory.onFileCompleted(writePlan);
     }
 
     private class PhysicalBundleRecords implements BundleRecords {
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingWritePlanFactory.java
 
b/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingWritePlanFactory.java
index 5adb313fe6..f2490fd046 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingWritePlanFactory.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingWritePlanFactory.java
@@ -36,4 +36,12 @@ public interface ShreddingWritePlanFactory {
     int inferBufferRowCount();
 
     ShreddingWritePlan createWritePlan(List<InternalRow> sampleRows);
+
+    /**
+     * Reports that a file using the given write plan has been closed 
successfully.
+     *
+     * <p>Implementations may use this callback to advance 
rolling-writer-scoped state. A factory is
+     * used sequentially by one rolling writer and is not thread-safe.
+     */
+    default void onFileCompleted(ShreddingWritePlan writePlan) {}
 }
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingWritePlanWriterFactory.java
 
b/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingWritePlanWriterFactory.java
index 5441eb4a97..9b4c6211f3 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingWritePlanWriterFactory.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/format/shredding/ShreddingWritePlanWriterFactory.java
@@ -55,6 +55,7 @@ public class ShreddingWritePlanWriterFactory implements 
FormatWriterFactory {
 
         return createWriterWithPlan(
                 delegate,
+                writePlanFactory,
                 out,
                 compression,
                 writePlanFactory.createWritePlan(Collections.emptyList()));
@@ -62,6 +63,7 @@ public class ShreddingWritePlanWriterFactory implements 
FormatWriterFactory {
 
     static FormatWriter createWriterWithPlan(
             SupportsShreddingWritePlan delegate,
+            ShreddingWritePlanFactory writePlanFactory,
             PositionOutputStream out,
             String compression,
             ShreddingWritePlan writePlan)
@@ -69,6 +71,7 @@ public class ShreddingWritePlanWriterFactory implements 
FormatWriterFactory {
         return new ShreddingFormatWriter(
                 delegate.createWithShreddingWritePlan(out, compression, 
writePlan),
                 delegate,
+                writePlanFactory,
                 writePlan,
                 compression);
     }
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/format/variant/VariantShreddingWritePlanFactory.java
 
b/paimon-common/src/main/java/org/apache/paimon/format/variant/VariantShreddingWritePlanFactory.java
index 75e4c77e68..357248c90a 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/format/variant/VariantShreddingWritePlanFactory.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/format/variant/VariantShreddingWritePlanFactory.java
@@ -22,6 +22,7 @@ import org.apache.paimon.CoreOptions;
 import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.data.shredding.ShreddingWritePlan;
 import org.apache.paimon.data.variant.InferVariantShreddingSchema;
+import org.apache.paimon.data.variant.VariantShreddingInferenceSession;
 import org.apache.paimon.data.variant.VariantShreddingWritePlan;
 import org.apache.paimon.format.shredding.ShreddingWritePlanFactory;
 import org.apache.paimon.options.Options;
@@ -30,6 +31,9 @@ import org.apache.paimon.types.DataType;
 import org.apache.paimon.types.RowType;
 import org.apache.paimon.types.VariantType;
 import org.apache.paimon.utils.JsonSerdeUtil;
+import org.apache.paimon.utils.Preconditions;
+
+import javax.annotation.Nullable;
 
 import java.util.List;
 
@@ -38,10 +42,25 @@ public class VariantShreddingWritePlanFactory implements 
ShreddingWritePlanFacto
 
     private final RowType rowType;
     private final Options options;
+    @Nullable private final VariantShreddingInferenceSession adaptiveSession;
+    @Nullable private ShreddingWritePlan pendingAdaptivePlan;
 
     public VariantShreddingWritePlanFactory(RowType rowType, Options options) {
         this.rowType = rowType;
         this.options = options;
+        if (isAdaptiveInference() && shouldInferWritePlan()) {
+            double admissionRatio =
+                    
options.get(CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO);
+            this.adaptiveSession =
+                    new VariantShreddingInferenceSession(
+                            createInferrer(),
+                            options.get(
+                                    
CoreOptions.VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW),
+                            admissionRatio,
+                            
options.get(CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO));
+        } else {
+            this.adaptiveSession = null;
+        }
     }
 
     @Override
@@ -69,6 +88,9 @@ public class VariantShreddingWritePlanFactory implements 
ShreddingWritePlanFacto
 
     @Override
     public int inferBufferRowCount() {
+        if (adaptiveSession != null && adaptiveSession.hasPrior()) {
+            return 
options.get(CoreOptions.VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW);
+        }
         return options.get(CoreOptions.VARIANT_SHREDDING_MAX_INFER_BUFFER_ROW);
     }
 
@@ -79,14 +101,41 @@ public class VariantShreddingWritePlanFactory implements 
ShreddingWritePlanFacto
                     rowType, configuredShreddingSchema());
         }
 
-        RowType physicalRowType = createInferrer().inferSchema(sampleRows);
-        return new VariantShreddingWritePlan(rowType, physicalRowType);
+        RowType physicalRowType;
+        if (adaptiveSession == null) {
+            physicalRowType = createInferrer().inferSchema(sampleRows);
+        } else {
+            physicalRowType = adaptiveSession.inferSchema(sampleRows);
+        }
+        VariantShreddingWritePlan writePlan =
+                new VariantShreddingWritePlan(rowType, physicalRowType);
+        if (adaptiveSession != null) {
+            pendingAdaptivePlan = writePlan;
+        }
+        return writePlan;
+    }
+
+    @Override
+    public void onFileCompleted(ShreddingWritePlan writePlan) {
+        if (adaptiveSession == null) {
+            return;
+        }
+        Preconditions.checkState(
+                writePlan == pendingAdaptivePlan,
+                "Completed Variant write plan does not match the pending 
inference.");
+        adaptiveSession.commitPendingInference();
+        pendingAdaptivePlan = null;
     }
 
     private boolean hasConfiguredShreddingSchema() {
         return options.contains(CoreOptions.VARIANT_SHREDDING_SCHEMA);
     }
 
+    private boolean isAdaptiveInference() {
+        return options.get(CoreOptions.VARIANT_SHREDDING_INFERENCE_MODE)
+                == CoreOptions.VariantShreddingInferenceMode.ADAPTIVE;
+    }
+
     private RowType configuredShreddingSchema() {
         String shreddingSchema = 
options.get(CoreOptions.VARIANT_SHREDDING_SCHEMA);
         return (RowType) JsonSerdeUtil.fromJson(shreddingSchema, 
DataType.class);
diff --git 
a/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanTest.java
 
b/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanTest.java
index 1652e81717..e10afa9652 100644
--- 
a/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanTest.java
+++ 
b/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanTest.java
@@ -82,7 +82,7 @@ class MapSharedShreddingWritePlanTest {
     }
 
     @Test
-    void testFactoryInfersColumnCountFromFirstRow() {
+    void testFactoryUsesMaxColumnCountForFirstFile() {
         RowType logicalType =
                 DataTypes.ROW(
                         DataTypes.FIELD(
@@ -90,8 +90,8 @@ class MapSharedShreddingWritePlanTest {
         MapSharedShreddingWritePlanFactory factory = 
createFactory(logicalType, 4);
 
         assertThat(factory.shouldCreateWritePlan()).isTrue();
-        assertThat(factory.shouldInferWritePlan()).isTrue();
-        assertThat(factory.inferBufferRowCount()).isEqualTo(1);
+        assertThat(factory.shouldInferWritePlan()).isFalse();
+        assertThat(factory.inferBufferRowCount()).isZero();
         assertThat(
                         factory.createWritePlan(
                                         Collections.singletonList(
@@ -100,26 +100,31 @@ class MapSharedShreddingWritePlanTest {
                                 .physicalRowType())
                 .isEqualTo(
                         MapSharedShreddingUtils.logicalToPhysicalSchema(
-                                logicalType, Collections.singletonMap("tags", 
3)));
+                                logicalType, Collections.singletonMap("tags", 
4)));
     }
 
     @Test
-    void testFactoryCapsInferredColumnCountAtMaxColumns() {
+    void testFactoryUsesCompletedFileStatisticsForNextFile() {
         RowType logicalType =
                 DataTypes.ROW(
                         DataTypes.FIELD(
                                 0, "tags", DataTypes.MAP(DataTypes.STRING(), 
DataTypes.INT())));
-        MapSharedShreddingWritePlanFactory factory = 
createFactory(logicalType, 2);
+        MapSharedShreddingWritePlanFactory factory = 
createFactory(logicalType, 8);
 
-        assertThat(
-                        factory.createWritePlan(
-                                        Collections.singletonList(
-                                                GenericRow.of(
-                                                        stringKeyMap("a", 1, 
"b", 2, "c", 3))))
-                                .physicalRowType())
+        ShreddingWritePlan firstPlan = 
factory.createWritePlan(Collections.emptyList());
+        firstPlan.toPhysicalRow(GenericRow.of(stringKeyMap("a", 1, "b", 
2))).getRow(0, 10);
+        firstPlan.toPhysicalRow(GenericRow.of(stringKeyMap("c", 3, "d", 4, 
"e", 5))).getRow(0, 10);
+        factory.onFileCompleted(firstPlan);
+
+        ShreddingWritePlan secondPlan = 
factory.createWritePlan(Collections.emptyList());
+        assertThat(firstPlan.physicalRowType())
                 .isEqualTo(
                         MapSharedShreddingUtils.logicalToPhysicalSchema(
-                                logicalType, Collections.singletonMap("tags", 
2)));
+                                logicalType, Collections.singletonMap("tags", 
8)));
+        assertThat(secondPlan.physicalRowType())
+                .isEqualTo(
+                        MapSharedShreddingUtils.logicalToPhysicalSchema(
+                                logicalType, Collections.singletonMap("tags", 
3)));
     }
 
     @Test
diff --git 
a/paimon-common/src/test/java/org/apache/paimon/data/variant/InferVariantShreddingSchemaTest.java
 
b/paimon-common/src/test/java/org/apache/paimon/data/variant/InferVariantShreddingSchemaTest.java
index 68787f8a45..ad541b0157 100644
--- 
a/paimon-common/src/test/java/org/apache/paimon/data/variant/InferVariantShreddingSchemaTest.java
+++ 
b/paimon-common/src/test/java/org/apache/paimon/data/variant/InferVariantShreddingSchemaTest.java
@@ -195,6 +195,103 @@ public class InferVariantShreddingSchemaTest {
                 .isEqualTo(variantShreddingSchema(DataTypes.VARIANT()));
     }
 
+    @Test
+    void testAdaptiveInferenceUsesAdmissionAndRetentionThresholds() {
+        RowType schema = RowType.of(new DataType[] {DataTypes.VARIANT()}, new 
String[] {"v"});
+        VariantShreddingInferenceSession session =
+                new VariantShreddingInferenceSession(
+                        new InferVariantShreddingSchema(schema, 300, 50, 0.4), 
10, 0.4, 0.2);
+
+        List<InternalRow> initialRows = new ArrayList<>();
+        for (int i = 0; i < 10; i++) {
+            initialRows.add(
+                    GenericRow.of(
+                            GenericVariant.fromJson(
+                                    i < 5 ? "{\"legacy\":\"v\",\"stable\":1}" 
: "{\"stable\":1}")));
+        }
+        RowType initialSchema = session.inferSchema(initialRows);
+        assertThat(initialSchema.getField("v").type())
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {DataTypes.STRING(), 
DataTypes.BIGINT()},
+                                        new String[] {"legacy", "stable"})));
+        session.commitPendingInference();
+
+        List<InternalRow> secondRows = new ArrayList<>();
+        for (int i = 0; i < 10; i++) {
+            secondRows.add(
+                    GenericRow.of(
+                            GenericVariant.fromJson(
+                                    i < 9
+                                            ? 
"{\"emerging\":true,\"stable\":2}"
+                                            : "{\"stable\":2}")));
+        }
+        RowType secondSchema = session.inferSchema(secondRows);
+        assertThat(secondSchema.getField("v").type())
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {
+                                            DataTypes.BOOLEAN(),
+                                            DataTypes.STRING(),
+                                            DataTypes.BIGINT()
+                                        },
+                                        new String[] {"emerging", "legacy", 
"stable"})));
+        session.commitPendingInference();
+
+        List<InternalRow> thirdRows = new ArrayList<>();
+        for (int i = 0; i < 10; i++) {
+            
thirdRows.add(GenericRow.of(GenericVariant.fromJson("{\"stable\":3}")));
+        }
+        RowType thirdSchema = session.inferSchema(thirdRows);
+        assertThat(thirdSchema.getField("v").type())
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {DataTypes.BOOLEAN(), 
DataTypes.BIGINT()},
+                                        new String[] {"emerging", "stable"})));
+    }
+
+    @Test
+    void testAdaptiveInferenceWidensScalarSelectedFromPriorEvidence() {
+        RowType schema =
+                RowType.of(
+                        new DataType[] {DataTypes.VARIANT(), 
DataTypes.VARIANT()},
+                        new String[] {"first", "second"});
+        VariantShreddingInferenceSession session =
+                new VariantShreddingInferenceSession(
+                        new InferVariantShreddingSchema(schema, 6, 50, 0.1), 
10, 0.1, 0.05);
+
+        RowType initialSchema =
+                session.inferSchema(
+                        Arrays.asList(
+                                GenericRow.of(
+                                        
GenericVariant.fromJson("{\"a\":1,\"b\":2}"),
+                                        
GenericVariant.fromJson("{\"historical\":12345}"))));
+        assertThat(initialSchema.getField("first").type())
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {DataTypes.BIGINT(), 
DataTypes.BIGINT()},
+                                        new String[] {"a", "b"})));
+        assertThat(initialSchema.getField("second").type())
+                .isEqualTo(variantShreddingSchema(DataTypes.VARIANT()));
+        session.commitPendingInference();
+
+        RowType adaptiveSchema =
+                session.inferSchema(
+                        
Arrays.asList(GenericRow.of(GenericVariant.fromJson("1"), null)));
+        assertThat(adaptiveSchema.getField("first").type())
+                .isEqualTo(variantShreddingSchema(DataTypes.BIGINT()));
+        assertThat(adaptiveSchema.getField("second").type())
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {DataTypes.BIGINT()},
+                                        new String[] {"historical"})));
+    }
+
     @Test
     void testInferSchemaWithDeepNesting() {
         // Schema: row<v: variant>
diff --git 
a/paimon-common/src/test/java/org/apache/paimon/format/shredding/ShreddingFormatWriterTest.java
 
b/paimon-common/src/test/java/org/apache/paimon/format/shredding/ShreddingFormatWriterTest.java
index 8345d2ae41..0bba5e23d8 100644
--- 
a/paimon-common/src/test/java/org/apache/paimon/format/shredding/ShreddingFormatWriterTest.java
+++ 
b/paimon-common/src/test/java/org/apache/paimon/format/shredding/ShreddingFormatWriterTest.java
@@ -28,6 +28,7 @@ import org.junit.jupiter.api.Test;
 
 import java.io.IOException;
 import java.util.Collections;
+import java.util.List;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -44,6 +45,7 @@ class ShreddingFormatWriterTest {
                         delegate,
                         new ThrowingMetadataFactory(failure),
                         NoOpWritePlan.INSTANCE,
+                        NoOpWritePlan.INSTANCE,
                         "none");
 
         assertThatThrownBy(writer::close).isSameAs(failure);
@@ -90,7 +92,7 @@ class ShreddingFormatWriterTest {
         }
     }
 
-    private enum NoOpWritePlan implements ShreddingWritePlan {
+    private enum NoOpWritePlan implements ShreddingWritePlan, 
ShreddingWritePlanFactory {
         INSTANCE;
 
         private final RowType rowType = new RowType(Collections.emptyList());
@@ -109,5 +111,25 @@ class ShreddingFormatWriterTest {
         public InternalRow toPhysicalRow(InternalRow row) {
             return row;
         }
+
+        @Override
+        public boolean shouldCreateWritePlan() {
+            return true;
+        }
+
+        @Override
+        public boolean shouldInferWritePlan() {
+            return false;
+        }
+
+        @Override
+        public int inferBufferRowCount() {
+            return 0;
+        }
+
+        @Override
+        public ShreddingWritePlan createWritePlan(List<InternalRow> 
sampleRows) {
+            return this;
+        }
     }
 }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java 
b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java
index 5644aab375..7c55360fef 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java
@@ -136,6 +136,7 @@ public class KeyValueFileWriterFactory {
     public RollingFileWriter<KeyValue, DataFileMeta> 
createRollingMergeTreeFileWriter(
             int level, FileSource fileSource) {
         WriteFormatKey key = new WriteFormatKey(level, false);
+        FormatWriterFactory writerFactory = 
formatContext.createWriterFactory(key);
         // Row limit applies to writes only; compaction output stays size-only.
         long targetFileRowNum =
                 fileSource == FileSource.COMPACT ? Long.MAX_VALUE : 
options.targetFileRowNum();
@@ -143,7 +144,11 @@ public class KeyValueFileWriterFactory {
                 () -> {
                     DataFilePathFactory pathFactory = 
formatContext.pathFactory(key);
                     return createDataFileWriter(
-                            pathFactory.newPath(), key, fileSource, 
pathFactory.isExternalPath());
+                            pathFactory.newPath(),
+                            key,
+                            fileSource,
+                            pathFactory.isExternalPath(),
+                            writerFactory);
                 },
                 suggestedFileSize,
                 targetFileRowNum);
@@ -151,6 +156,7 @@ public class KeyValueFileWriterFactory {
 
     public RollingFileWriter<KeyValue, DataFileMeta> 
createRollingChangelogFileWriter(int level) {
         WriteFormatKey key = new WriteFormatKey(level, true);
+        FormatWriterFactory writerFactory = 
formatContext.createWriterFactory(key);
         return new RollingFileWriterImpl<>(
                 () -> {
                     DataFilePathFactory pathFactory = 
formatContext.pathFactory(key);
@@ -158,7 +164,8 @@ public class KeyValueFileWriterFactory {
                             pathFactory.newChangelogPath(),
                             key,
                             FileSource.APPEND,
-                            pathFactory.isExternalPath());
+                            pathFactory.isExternalPath(),
+                            writerFactory);
                 },
                 suggestedFileSize,
                 Long.MAX_VALUE);
@@ -166,21 +173,28 @@ public class KeyValueFileWriterFactory {
 
     public RollingFileWriter<KeyValue, DataFileMeta> 
createRollingClusteringFileWriter() {
         WriteFormatKey key = new WriteFormatKey(1, false);
+        FormatWriterFactory writerFactory = 
formatContext.createWriterFactory(key);
         return new RollingFileWriterImpl<>(
                 () -> {
                     DataFilePathFactory pathFactory = 
formatContext.pathFactory(key);
                     return createKvSeparatedFileWriter(
-                            pathFactory.newPath(), key, 
pathFactory.isExternalPath());
+                            pathFactory.newPath(),
+                            key,
+                            pathFactory.isExternalPath(),
+                            writerFactory);
                 },
                 suggestedFileSize,
                 Long.MAX_VALUE);
     }
 
     private KeyValueClusteringFileWriter createKvSeparatedFileWriter(
-            Path path, WriteFormatKey key, boolean isExternalPath) {
+            Path path,
+            WriteFormatKey key,
+            boolean isExternalPath,
+            FormatWriterFactory writerFactory) {
         return new KeyValueClusteringFileWriter(
                 fileIO,
-                formatContext.fileWriterContext(key),
+                formatContext.fileWriterContext(key, writerFactory),
                 path,
                 keyType,
                 valueType,
@@ -193,7 +207,11 @@ public class KeyValueFileWriterFactory {
     }
 
     private KeyValueDataFileWriter createDataFileWriter(
-            Path path, WriteFormatKey key, FileSource fileSource, boolean 
isExternalPath) {
+            Path path,
+            WriteFormatKey key,
+            FileSource fileSource,
+            boolean isExternalPath,
+            FormatWriterFactory writerFactory) {
         // Changelog is sequentially consumed, file index is unnecessary.
         FileIndexOptions indexOptions = key.isChangelog ? new 
FileIndexOptions() : fileIndexOptions;
         Set<String> dataFileManagedBlobFields =
@@ -201,7 +219,7 @@ public class KeyValueFileWriterFactory {
         return formatContext.thinModeEnabled
                 ? new KeyValueThinDataFileWriterImpl(
                         fileIO,
-                        formatContext.fileWriterContext(key),
+                        formatContext.fileWriterContext(key, writerFactory),
                         path,
                         new KeyValueThinSerializer(keyType, valueType)::toRow,
                         keyType,
@@ -215,7 +233,7 @@ public class KeyValueFileWriterFactory {
                         dataFileManagedBlobFields)
                 : new KeyValueDataFileWriterImpl(
                         fileIO,
-                        formatContext.fileWriterContext(key),
+                        formatContext.fileWriterContext(key, writerFactory),
                         path,
                         new KeyValueSerializer(keyType, valueType)::toRow,
                         keyType,
@@ -322,7 +340,6 @@ public class KeyValueFileWriterFactory {
         private final Map<String, SimpleColStatsCollector.Factory[]> 
statsMode2AvroStats;
         private final Map<String, DataFilePathFactory> format2PathFactory;
         private final Map<String, FileFormat> formatFactory;
-        private final Map<String, FormatWriterFactory> format2WriterFactory;
 
         private final BinaryRow partition;
         private final int bucket;
@@ -388,7 +405,6 @@ public class KeyValueFileWriterFactory {
             this.formatStats2Extractor = new HashMap<>();
             this.statsMode2AvroStats = new HashMap<>();
             this.format2PathFactory = new HashMap<>();
-            this.format2WriterFactory = new HashMap<>();
             this.formatFactory = new HashMap<>();
         }
 
@@ -407,9 +423,10 @@ public class KeyValueFileWriterFactory {
             return true;
         }
 
-        private FileWriterContext fileWriterContext(WriteFormatKey key) {
+        private FileWriterContext fileWriterContext(
+                WriteFormatKey key, FormatWriterFactory writerFactory) {
             return new FileWriterContext(
-                    writerFactory(key), statsProducer(key), 
key2Compress.apply(key));
+                    writerFactory, statsProducer(key), 
key2Compress.apply(key));
         }
 
         private SimpleStatsProducer statsProducer(WriteFormatKey key) {
@@ -462,10 +479,8 @@ public class KeyValueFileWriterFactory {
                                     .createDataFilePathFactory(partition, 
bucket));
         }
 
-        private FormatWriterFactory writerFactory(WriteFormatKey key) {
-            return format2WriterFactory.computeIfAbsent(
-                    key2Format.apply(key),
-                    format -> 
fileFormat(format).createWriterFactory(writeRowType));
+        private FormatWriterFactory createWriterFactory(WriteFormatKey key) {
+            return 
fileFormat(key2Format.apply(key)).createWriterFactory(writeRowType);
         }
 
         private FileFormat fileFormat(String format) {
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/io/RowDataRollingFileWriter.java 
b/paimon-core/src/main/java/org/apache/paimon/io/RowDataRollingFileWriter.java
index e9874ef9b9..9b35c002aa 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/io/RowDataRollingFileWriter.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/io/RowDataRollingFileWriter.java
@@ -21,6 +21,7 @@ package org.apache.paimon.io;
 import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.fileindex.FileIndexOptions;
 import org.apache.paimon.format.FileFormat;
+import org.apache.paimon.format.FormatWriterFactory;
 import org.apache.paimon.fs.FileIO;
 import org.apache.paimon.fs.Path;
 import org.apache.paimon.manifest.FileSource;
@@ -54,28 +55,41 @@ public class RowDataRollingFileWriter extends 
RollingFileWriterImpl<InternalRow,
             @Nullable FileFormat rowSidecarFormat,
             long targetFileRowNum) {
         super(
-                () -> {
-                    Path dataPath = pathFactory.newPath();
-                    Path rowSidecarPath =
-                            rowSidecarFormat == null
-                                    ? null
-                                    : new Path(dataPath.getParent(), 
dataPath.getName() + ".row");
-                    return new RowDataFileWriter(
-                            fileIO,
-                            RollingFileWriter.createFileWriterContext(
-                                    fileFormat, writeSchema, statsCollectors, 
fileCompression),
-                            dataPath,
-                            writeSchema,
-                            schemaId,
-                            seqNumCounterSupplier,
-                            fileIndexOptions,
-                            fileSource,
-                            asyncFileWrite,
-                            statsDenseStore,
-                            pathFactory.isExternalPath(),
-                            writeCols,
-                            rowSidecarFormat,
-                            rowSidecarPath);
+                new Supplier<RowDataFileWriter>() {
+
+                    private final FormatWriterFactory formatWriterFactory =
+                            fileFormat.createWriterFactory(writeSchema);
+
+                    @Override
+                    public RowDataFileWriter get() {
+                        Path dataPath = pathFactory.newPath();
+                        Path rowSidecarPath =
+                                rowSidecarFormat == null
+                                        ? null
+                                        : new Path(
+                                                dataPath.getParent(), 
dataPath.getName() + ".row");
+                        FileWriterContext writerContext =
+                                new FileWriterContext(
+                                        formatWriterFactory,
+                                        RollingFileWriter.createStatsProducer(
+                                                fileFormat, writeSchema, 
statsCollectors),
+                                        fileCompression);
+                        return new RowDataFileWriter(
+                                fileIO,
+                                writerContext,
+                                dataPath,
+                                writeSchema,
+                                schemaId,
+                                seqNumCounterSupplier,
+                                fileIndexOptions,
+                                fileSource,
+                                asyncFileWrite,
+                                statsDenseStore,
+                                pathFactory.isExternalPath(),
+                                writeCols,
+                                rowSidecarFormat,
+                                rowSidecarPath);
+                    }
                 },
                 targetFileSize,
                 targetFileRowNum);
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java 
b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java
index 1c6ba3a0fb..cdd2555c36 100644
--- a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java
+++ b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java
@@ -367,6 +367,7 @@ public class SchemaValidation {
         validateMergeFunctionFactory(schema);
 
         validateMapStorageLayout(schema, options);
+        validateVariantShreddingInferenceOptions(options);
 
         validateFileIndex(schema);
 
@@ -695,6 +696,44 @@ public class SchemaValidation {
         }
     }
 
+    private static void validateVariantShreddingInferenceOptions(CoreOptions 
options) {
+        int coldSampleRows =
+                
options.toConfiguration().get(CoreOptions.VARIANT_SHREDDING_MAX_INFER_BUFFER_ROW);
+        checkArgument(
+                coldSampleRows > 0,
+                "%s must be positive.",
+                CoreOptions.VARIANT_SHREDDING_MAX_INFER_BUFFER_ROW.key());
+
+        double admissionRatio =
+                options.toConfiguration()
+                        
.get(CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO);
+        checkArgument(
+                admissionRatio >= 0 && admissionRatio <= 1,
+                "%s must be between 0 and 1.",
+                
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key());
+
+        if 
(options.toConfiguration().get(CoreOptions.VARIANT_SHREDDING_INFERENCE_MODE)
+                != CoreOptions.VariantShreddingInferenceMode.ADAPTIVE) {
+            return;
+        }
+
+        int warmSampleRows =
+                options.toConfiguration()
+                        
.get(CoreOptions.VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW);
+        checkArgument(
+                warmSampleRows > 0,
+                "%s must be positive.",
+                
CoreOptions.VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW.key());
+        double retentionRatio =
+                options.toConfiguration()
+                        
.get(CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO);
+        checkArgument(
+                retentionRatio >= 0 && retentionRatio <= admissionRatio,
+                "%s must be between 0 and %s.",
+                CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO.key(),
+                
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key());
+    }
+
     private static void validateUnsupportedTypesWithMapSharedShredding(
             TableSchema schema, CoreOptions options) {
         RowType rowType = new RowType(schema.fields());
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java 
b/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java
index 45f2a8a4b1..ef56732447 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java
@@ -509,7 +509,7 @@ public class AppendOnlyWriterTest {
                                 nameToId("a", 0, "b", 1),
                                 fieldToColumns(0, columns(0), 1, columns(0)),
                                 overflowFields(),
-                                1,
+                                3,
                                 1));
     }
 
@@ -528,7 +528,7 @@ public class AppendOnlyWriterTest {
 
         assertThat(readSharedShreddingFieldMeta(context, file, "tags"))
                 .isEqualTo(
-                        sharedShreddingMeta(nameToId(), fieldToColumns(), 
overflowFields(), 1, 0));
+                        sharedShreddingMeta(nameToId(), fieldToColumns(), 
overflowFields(), 3, 0));
     }
 
     @ParameterizedTest(name = "{0}")
@@ -557,17 +557,18 @@ public class AppendOnlyWriterTest {
 
         assertThat(readSharedShreddingFieldMeta(context, nullFile, "tags"))
                 .isEqualTo(
-                        sharedShreddingMeta(nameToId(), fieldToColumns(), 
overflowFields(), 1, 0));
+                        sharedShreddingMeta(nameToId(), fieldToColumns(), 
overflowFields(), 3, 0));
         assertThat(readSharedShreddingFieldMeta(context, emptyFile, "tags"))
                 .isEqualTo(
-                        sharedShreddingMeta(nameToId(), fieldToColumns(), 
overflowFields(), 1, 0));
+                        sharedShreddingMeta(nameToId(), fieldToColumns(), 
overflowFields(), 3, 0));
         assertThat(readSharedShreddingFieldMeta(context, valuesFile, "tags"))
                 .isEqualTo(
                         sharedShreddingMeta(
                                 nameToId("a", 0, "b", 1, "c", 2, "d", 3),
-                                fieldToColumns(0, columns(0), 1, columns(0), 
2, columns(0)),
-                                overflowFields(3),
-                                1,
+                                fieldToColumns(
+                                        0, columns(0), 1, columns(0), 2, 
columns(0), 3, columns(1)),
+                                overflowFields(),
+                                3,
                                 2));
     }
 
@@ -613,9 +614,9 @@ public class AppendOnlyWriterTest {
                 .isEqualTo(
                         sharedShreddingMeta(
                                 nameToId("a", 0, "b", 1, "c", 2),
-                                fieldToColumns(0, columns(0), 1, columns(1)),
-                                overflowFields(2),
-                                2,
+                                fieldToColumns(0, columns(0), 1, columns(1), 
2, columns(2)),
+                                overflowFields(),
+                                3,
                                 3));
     }
 
@@ -658,7 +659,7 @@ public class AppendOnlyWriterTest {
                                 nameToId("a", 0),
                                 fieldToColumns(0, columns(0)),
                                 overflowFields(),
-                                1,
+                                3,
                                 1));
     }
 
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java 
b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java
index ab94f33b3a..36a085c6a8 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java
@@ -120,6 +120,83 @@ class SchemaValidationTest {
                 .hasMessageContaining("target-file-row-num should be at least 
1");
     }
 
+    @Test
+    public void testVariantShreddingInferenceOptions() {
+        Map<String, String> invalidColdSampleOptions = new HashMap<>();
+        
invalidColdSampleOptions.put(CoreOptions.VARIANT_SHREDDING_MAX_INFER_BUFFER_ROW.key(),
 "0");
+        assertThatThrownBy(() -> 
validateTableSchemaExec(invalidColdSampleOptions))
+                .hasMessageContaining(
+                        
CoreOptions.VARIANT_SHREDDING_MAX_INFER_BUFFER_ROW.key()
+                                + " must be positive");
+
+        Map<String, String> invalidAdmissionOptions = new HashMap<>();
+        invalidAdmissionOptions.put(
+                
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key(), "1.1");
+        assertThatThrownBy(() -> 
validateTableSchemaExec(invalidAdmissionOptions))
+                .hasMessageContaining(
+                        
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key()
+                                + " must be between 0 and 1");
+
+        Map<String, String> negativeAdmissionOptions = new HashMap<>();
+        negativeAdmissionOptions.put(
+                
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key(), "-0.1");
+        assertThatThrownBy(() -> 
validateTableSchemaExec(negativeAdmissionOptions))
+                .hasMessageContaining(
+                        
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key()
+                                + " must be between 0 and 1");
+
+        Map<String, String> invalidWarmSampleOptions = new HashMap<>();
+        invalidWarmSampleOptions.put(
+                CoreOptions.VARIANT_SHREDDING_INFERENCE_MODE.key(), 
"adaptive");
+        invalidWarmSampleOptions.put(
+                
CoreOptions.VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW.key(), "0");
+        assertThatThrownBy(() -> 
validateTableSchemaExec(invalidWarmSampleOptions))
+                .hasMessageContaining(
+                        
CoreOptions.VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW.key()
+                                + " must be positive");
+
+        Map<String, String> invalidRetentionOptions = new HashMap<>();
+        
invalidRetentionOptions.put(CoreOptions.VARIANT_SHREDDING_INFERENCE_MODE.key(), 
"adaptive");
+        invalidRetentionOptions.put(
+                
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key(), "0.2");
+        invalidRetentionOptions.put(
+                CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO.key(), 
"0.3");
+        assertThatThrownBy(() -> 
validateTableSchemaExec(invalidRetentionOptions))
+                .hasMessageContaining(
+                        
CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO.key()
+                                + " must be between 0 and "
+                                + 
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key());
+
+        Map<String, String> negativeRetentionOptions = new HashMap<>();
+        negativeRetentionOptions.put(
+                CoreOptions.VARIANT_SHREDDING_INFERENCE_MODE.key(), 
"adaptive");
+        negativeRetentionOptions.put(
+                CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO.key(), 
"-0.1");
+        assertThatThrownBy(() -> 
validateTableSchemaExec(negativeRetentionOptions))
+                .hasMessageContaining(
+                        
CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO.key()
+                                + " must be between 0 and "
+                                + 
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key());
+
+        Map<String, String> perFileOptions = new HashMap<>();
+        perFileOptions.put(CoreOptions.VARIANT_SHREDDING_INFERENCE_MODE.key(), 
"per-file");
+        
perFileOptions.put(CoreOptions.VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW.key(),
 "0");
+        
perFileOptions.put(CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO.key(),
 "2");
+        assertThatCode(() -> 
validateTableSchemaExec(perFileOptions)).doesNotThrowAnyException();
+
+        Map<String, String> validAdaptiveOptions = new HashMap<>();
+        
validAdaptiveOptions.put(CoreOptions.VARIANT_SHREDDING_INFERENCE_MODE.key(), 
"adaptive");
+        
validAdaptiveOptions.put(CoreOptions.VARIANT_SHREDDING_MAX_INFER_BUFFER_ROW.key(),
 "10");
+        validAdaptiveOptions.put(
+                
CoreOptions.VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW.key(), "2");
+        validAdaptiveOptions.put(
+                
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key(), "0.2");
+        validAdaptiveOptions.put(
+                CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO.key(), 
"0.1");
+        assertThatCode(() -> validateTableSchemaExec(validAdaptiveOptions))
+                .doesNotThrowAnyException();
+    }
+
     @Test
     public void testFromSnapshotConflict() {
         String timestampString =
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/MapSharedShreddingTableTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/MapSharedShreddingTableTest.java
index 666654231b..fe06144044 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/MapSharedShreddingTableTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/MapSharedShreddingTableTest.java
@@ -218,7 +218,7 @@ public class MapSharedShreddingTableTest extends 
TableTestBase {
         
assertThat(files.get(0).dataFile.fileSource()).hasValue(FileSource.COMPACT);
         MapSharedShreddingFieldMeta compactedMeta =
                 readSharedShreddingFieldMeta(fileStoreTable, files.get(0), 
"metrics");
-        assertThat(compactedMeta.numColumns()).isEqualTo(2);
+        assertThat(compactedMeta.numColumns()).isEqualTo(64);
         assertThat(compactedMeta.maxRowWidth()).isEqualTo(4);
         assertThat(compactedMeta.nameToId()).containsOnlyKeys("a", "b", "c", 
"d", "e", "f", "g");
 
@@ -381,14 +381,21 @@ public class MapSharedShreddingTableTest extends 
TableTestBase {
 
     @ParameterizedTest
     @ValueSource(strings = {"orc", "parquet"})
-    public void testInferColumnCountFromFirstRowOfEachFile(String format) 
throws Exception {
+    public void testAdaptColumnCountAcrossRollingFiles(String format) throws 
Exception {
         Table table = createTableWithBucket(format, 8, "1", "metrics");
+        catalog.alterTable(
+                identifier(format),
+                Collections.singletonList(
+                        
SchemaChange.setOption(CoreOptions.TARGET_FILE_ROW_NUM.key(), "2")),
+                false);
+        table = catalog.getTable(identifier(format));
 
         write(
                 table,
                 GenericRow.of(1, mapOf("a", 11L, "b", 12L)),
-                GenericRow.of(2, mapOf("c", 21L, "d", 22L, "e", 23L)));
-        write(table, GenericRow.of(3, mapOf("f", 31L)));
+                GenericRow.of(2, mapOf("c", 21L, "d", 22L, "e", 23L)),
+                GenericRow.of(3, mapOf("f", 31L)),
+                GenericRow.of(4, mapOf("g", 41L, "h", 42L)));
 
         FileStoreTable fileStoreTable = (FileStoreTable) table;
         List<DataFileWithSplit> files = currentDataFiles(fileStoreTable);
@@ -397,13 +404,13 @@ public class MapSharedShreddingTableTest extends 
TableTestBase {
 
         MapSharedShreddingFieldMeta firstFileMeta =
                 readSharedShreddingFieldMeta(fileStoreTable, files.get(0), 
"metrics");
-        assertThat(firstFileMeta.numColumns()).isEqualTo(2);
+        assertThat(firstFileMeta.numColumns()).isEqualTo(8);
         assertThat(firstFileMeta.maxRowWidth()).isEqualTo(3);
 
         MapSharedShreddingFieldMeta secondFileMeta =
                 readSharedShreddingFieldMeta(fileStoreTable, files.get(1), 
"metrics");
-        assertThat(secondFileMeta.numColumns()).isEqualTo(1);
-        assertThat(secondFileMeta.maxRowWidth()).isEqualTo(1);
+        assertThat(secondFileMeta.numColumns()).isEqualTo(3);
+        assertThat(secondFileMeta.maxRowWidth()).isEqualTo(2);
 
         Map<Integer, Map<String, Long>> actual = new LinkedHashMap<>();
         for (InternalRow row : read(table)) {
@@ -413,7 +420,8 @@ public class MapSharedShreddingTableTest extends 
TableTestBase {
         assertThat(actual)
                 .containsEntry(1, javaMapOf("a", 11L, "b", 12L))
                 .containsEntry(2, javaMapOf("c", 21L, "d", 22L, "e", 23L))
-                .containsEntry(3, javaMapOf("f", 31L));
+                .containsEntry(3, javaMapOf("f", 31L))
+                .containsEntry(4, javaMapOf("g", 41L, "h", 42L));
     }
 
     @ParameterizedTest
@@ -690,8 +698,8 @@ public class MapSharedShreddingTableTest extends 
TableTestBase {
 
     @ParameterizedTest
     @CsvSource({"orc,false", "orc,true", "parquet,false", "parquet,true"})
-    public void testPrimaryKeyInfersColumnCountPerFile(String format, boolean 
thinMode)
-            throws Exception {
+    public void testPrimaryKeyStartsNewRollingWriterWithMaxColumnCount(
+            String format, boolean thinMode) throws Exception {
         Table table = createPrimaryKeyTable(format, thinMode, 8);
 
         write(table, GenericRow.of(1, mapOf("a", 11L, "b", 12L)));
@@ -704,12 +712,12 @@ public class MapSharedShreddingTableTest extends 
TableTestBase {
 
         MapSharedShreddingFieldMeta firstFileMeta =
                 readSharedShreddingFieldMeta(fileStoreTable, files.get(0), 
"metrics");
-        assertThat(firstFileMeta.numColumns()).isEqualTo(2);
+        assertThat(firstFileMeta.numColumns()).isEqualTo(8);
         assertThat(firstFileMeta.maxRowWidth()).isEqualTo(2);
 
         MapSharedShreddingFieldMeta secondFileMeta =
                 readSharedShreddingFieldMeta(fileStoreTable, files.get(1), 
"metrics");
-        assertThat(secondFileMeta.numColumns()).isEqualTo(1);
+        assertThat(secondFileMeta.numColumns()).isEqualTo(8);
         assertThat(secondFileMeta.maxRowWidth()).isEqualTo(1);
 
         Map<Integer, Map<String, Long>> actual = new LinkedHashMap<>();
@@ -754,7 +762,7 @@ public class MapSharedShreddingTableTest extends 
TableTestBase {
 
         MapSharedShreddingFieldMeta metricsMeta =
                 readSharedShreddingFieldMeta(fileStoreTable, files.get(1), 
"metrics");
-        assertThat(metricsMeta.numColumns()).isEqualTo(1);
+        assertThat(metricsMeta.numColumns()).isEqualTo(3);
         assertThat(metricsMeta.maxRowWidth()).isEqualTo(1);
 
         Map<Integer, List<Map<String, Long>>> actual = new LinkedHashMap<>();
diff --git 
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/writer/InferVariantShreddingWriteTest.java
 
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/writer/InferVariantShreddingWriteTest.java
index 3e6b30df2f..2d25413616 100644
--- 
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/writer/InferVariantShreddingWriteTest.java
+++ 
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/writer/InferVariantShreddingWriteTest.java
@@ -121,6 +121,243 @@ public class InferVariantShreddingWriteTest {
         assertThat(result2.get(2)).isEqualTo(GenericRow.of(GenericRow.of(35)));
     }
 
+    @Test
+    public void testAdaptiveInferenceAcrossFiles() throws Exception {
+        ParquetFileFormat format = createFormat(adaptiveOptions(10, 10, 0.4, 
0.2));
+        RowType writeType = DataTypes.ROW(DataTypes.FIELD(0, "payload", 
DataTypes.VARIANT()));
+        FormatWriterFactory factory = format.createWriterFactory(writeType);
+
+        List<InternalRow> firstRows = new ArrayList<>();
+        List<String> firstJson = new ArrayList<>();
+        for (int i = 0; i < 10; i++) {
+            String json =
+                    i < 5
+                            ? "{\"legacy\":\"value\",\"stable\":" + i + "}"
+                            : "{\"stable\":" + i + "}";
+            firstRows.add(GenericRow.of(GenericVariant.fromJson(json)));
+            firstJson.add(json);
+        }
+        Path firstFile = file;
+        writeRows(factory, firstFile, firstRows.toArray(new InternalRow[0]));
+
+        List<InternalRow> secondRows = new ArrayList<>();
+        List<String> secondJson = new ArrayList<>();
+        for (int i = 0; i < 10; i++) {
+            String json =
+                    i < 9 ? "{\"emerging\":true,\"stable\":" + i + "}" : 
"{\"stable\":" + i + "}";
+            secondRows.add(GenericRow.of(GenericVariant.fromJson(json)));
+            secondJson.add(json);
+        }
+        Path secondFile = newFile();
+        writeRows(factory, secondFile, secondRows.toArray(new InternalRow[0]));
+
+        List<InternalRow> thirdRows = new ArrayList<>();
+        List<String> thirdJson = new ArrayList<>();
+        for (int i = 0; i < 10; i++) {
+            String json = "{\"stable\":" + i + "}";
+            thirdRows.add(GenericRow.of(GenericVariant.fromJson(json)));
+            thirdJson.add(json);
+        }
+        Path thirdFile = newFile();
+        writeRows(factory, thirdFile, thirdRows.toArray(new InternalRow[0]));
+
+        assertThat(readVariantFileType(firstFile, "payload"))
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {DataTypes.STRING(), 
DataTypes.BIGINT()},
+                                        new String[] {"legacy", "stable"})));
+        assertThat(readVariantFileType(secondFile, "payload"))
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {
+                                            DataTypes.BOOLEAN(),
+                                            DataTypes.STRING(),
+                                            DataTypes.BIGINT()
+                                        },
+                                        new String[] {"emerging", "legacy", 
"stable"})));
+        assertThat(readVariantFileType(thirdFile, "payload"))
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {DataTypes.BOOLEAN(), 
DataTypes.BIGINT()},
+                                        new String[] {"emerging", "stable"})));
+
+        assertThat(readVariantJson(format, writeType, firstFile, 0))
+                .containsExactlyElementsOf(firstJson);
+        assertThat(readVariantJson(format, writeType, secondFile, 0))
+                .containsExactlyElementsOf(secondJson);
+        assertThat(readVariantJson(format, writeType, thirdFile, 0))
+                .containsExactlyElementsOf(thirdJson);
+    }
+
+    @Test
+    public void testAdaptiveInferenceOnShortRolledFile() throws Exception {
+        ParquetFileFormat format = createFormat(adaptiveOptions(4, 4, 0.4, 
0.2));
+        RowType writeType = DataTypes.ROW(DataTypes.FIELD(0, "payload", 
DataTypes.VARIANT()));
+        FormatWriterFactory factory = format.createWriterFactory(writeType);
+
+        Path firstFile = file;
+        writeRows(
+                factory,
+                firstFile,
+                
GenericRow.of(GenericVariant.fromJson("{\"legacy\":\"a\",\"stable\":1}")),
+                
GenericRow.of(GenericVariant.fromJson("{\"legacy\":\"b\",\"stable\":2}")),
+                
GenericRow.of(GenericVariant.fromJson("{\"legacy\":\"c\",\"stable\":3}")),
+                
GenericRow.of(GenericVariant.fromJson("{\"legacy\":\"d\",\"stable\":4}")));
+
+        Path secondFile = newFile();
+        writeRows(
+                factory,
+                secondFile,
+                
GenericRow.of(GenericVariant.fromJson("{\"emerging\":true,\"stable\":5}")),
+                
GenericRow.of(GenericVariant.fromJson("{\"emerging\":false,\"stable\":6}")),
+                
GenericRow.of(GenericVariant.fromJson("{\"emerging\":true,\"stable\":7}")));
+
+        assertThat(readVariantFileType(secondFile, "payload"))
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {
+                                            DataTypes.BOOLEAN(),
+                                            DataTypes.STRING(),
+                                            DataTypes.BIGINT()
+                                        },
+                                        new String[] {"emerging", "legacy", 
"stable"})));
+        assertThat(readVariantJson(format, writeType, secondFile, 0))
+                .containsExactly(
+                        "{\"emerging\":true,\"stable\":5}",
+                        "{\"emerging\":false,\"stable\":6}",
+                        "{\"emerging\":true,\"stable\":7}");
+    }
+
+    @Test
+    public void testAdaptiveInferenceWithMultipleVariantFields() throws 
Exception {
+        ParquetFileFormat format = createFormat(adaptiveOptions(2, 2, 0.4, 
0.2));
+        RowType writeType =
+                DataTypes.ROW(
+                        DataTypes.FIELD(0, "left_payload", 
DataTypes.VARIANT()),
+                        DataTypes.FIELD(1, "right_payload", 
DataTypes.VARIANT()));
+        FormatWriterFactory factory = format.createWriterFactory(writeType);
+
+        Path firstFile = file;
+        writeRows(
+                factory,
+                firstFile,
+                GenericRow.of(
+                        
GenericVariant.fromJson("{\"legacy\":\"a\",\"stable\":1}"),
+                        
GenericVariant.fromJson("{\"sparse\":\"x\",\"stable\":\"a\"}")),
+                GenericRow.of(
+                        
GenericVariant.fromJson("{\"legacy\":\"b\",\"stable\":2}"),
+                        GenericVariant.fromJson("{\"stable\":\"b\"}")));
+
+        Path secondFile = newFile();
+        writeRows(
+                factory,
+                secondFile,
+                GenericRow.of(
+                        
GenericVariant.fromJson("{\"emerging\":true,\"stable\":3}"),
+                        GenericVariant.fromJson("{\"stable\":\"c\"}")),
+                GenericRow.of(
+                        
GenericVariant.fromJson("{\"emerging\":false,\"stable\":4}"),
+                        
GenericVariant.fromJson("{\"emerging\":true,\"stable\":\"d\"}")));
+
+        assertThat(readVariantFileType(firstFile, "left_payload"))
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {DataTypes.STRING(), 
DataTypes.BIGINT()},
+                                        new String[] {"legacy", "stable"})));
+        assertThat(readVariantFileType(secondFile, "left_payload"))
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {
+                                            DataTypes.BOOLEAN(),
+                                            DataTypes.STRING(),
+                                            DataTypes.BIGINT()
+                                        },
+                                        new String[] {"emerging", "legacy", 
"stable"})));
+
+        RowType rightSchema =
+                variantShreddingSchema(
+                        RowType.of(
+                                new DataType[] {DataTypes.STRING(), 
DataTypes.STRING()},
+                                new String[] {"sparse", "stable"}));
+        assertThat(readVariantFileType(firstFile, 
"right_payload")).isEqualTo(rightSchema);
+        assertThat(readVariantFileType(secondFile, 
"right_payload")).isEqualTo(rightSchema);
+
+        List<InternalRow> result = readRows(format, writeType, secondFile);
+        assertThat(result.get(0).getVariant(0).toJson())
+                .isEqualTo("{\"emerging\":true,\"stable\":3}");
+        
assertThat(result.get(0).getVariant(1).toJson()).isEqualTo("{\"stable\":\"c\"}");
+        assertThat(result.get(1).getVariant(0).toJson())
+                .isEqualTo("{\"emerging\":false,\"stable\":4}");
+        assertThat(result.get(1).getVariant(1).toJson())
+                .isEqualTo("{\"emerging\":true,\"stable\":\"d\"}");
+    }
+
+    @Test
+    public void testAdaptiveInferenceWithNestedVariant() throws Exception {
+        RowType nestedType =
+                DataTypes.ROW(
+                        DataTypes.FIELD(1, "label", DataTypes.STRING()),
+                        DataTypes.FIELD(2, "payload", DataTypes.VARIANT()));
+        RowType writeType = DataTypes.ROW(DataTypes.FIELD(0, "nested", 
nestedType));
+        ParquetFileFormat format = createFormat(adaptiveOptions(2, 2, 0.4, 
0.2));
+        FormatWriterFactory factory = format.createWriterFactory(writeType);
+
+        Path firstFile = file;
+        writeRows(
+                factory,
+                firstFile,
+                GenericRow.of(
+                        GenericRow.of(
+                                BinaryString.fromString("first"),
+                                
GenericVariant.fromJson("{\"legacy\":\"a\",\"stable\":1}"))),
+                GenericRow.of(
+                        GenericRow.of(
+                                BinaryString.fromString("second"),
+                                
GenericVariant.fromJson("{\"legacy\":\"b\",\"stable\":2}"))));
+
+        Path secondFile = newFile();
+        writeRows(
+                factory,
+                secondFile,
+                GenericRow.of(
+                        GenericRow.of(
+                                BinaryString.fromString("third"),
+                                
GenericVariant.fromJson("{\"emerging\":true,\"stable\":3}"))),
+                GenericRow.of(
+                        GenericRow.of(
+                                BinaryString.fromString("fourth"),
+                                
GenericVariant.fromJson("{\"emerging\":false,\"stable\":4}"))));
+
+        assertThat(readVariantFileType(firstFile, "nested", "payload"))
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {DataTypes.STRING(), 
DataTypes.BIGINT()},
+                                        new String[] {"legacy", "stable"})));
+        assertThat(readVariantFileType(secondFile, "nested", "payload"))
+                .isEqualTo(
+                        variantShreddingSchema(
+                                RowType.of(
+                                        new DataType[] {
+                                            DataTypes.BOOLEAN(),
+                                            DataTypes.STRING(),
+                                            DataTypes.BIGINT()
+                                        },
+                                        new String[] {"emerging", "legacy", 
"stable"})));
+
+        List<InternalRow> result = readRows(format, writeType, secondFile);
+        assertThat(result.get(0).getRow(0, 2).getVariant(1).toJson())
+                .isEqualTo("{\"emerging\":true,\"stable\":3}");
+        assertThat(result.get(1).getRow(0, 2).getVariant(1).toJson())
+                .isEqualTo("{\"emerging\":false,\"stable\":4}");
+    }
+
     @Test
     public void testInferSchemaWithArray() throws Exception {
         ParquetFileFormat format = createFormat();
@@ -513,26 +750,73 @@ public class InferVariantShreddingWriteTest {
         return new ParquetFileFormat(new 
FileFormatFactory.FormatContext(options, 1024, 1024));
     }
 
+    private Options adaptiveOptions(
+            int initialSampleRows,
+            int adaptiveSampleRows,
+            double admissionRatio,
+            double retentionRatio) {
+        Options options = defaultOptions();
+        options.set(CoreOptions.VARIANT_SHREDDING_INFERENCE_MODE.key(), 
"adaptive");
+        options.set(
+                CoreOptions.VARIANT_SHREDDING_MAX_INFER_BUFFER_ROW.key(),
+                String.valueOf(initialSampleRows));
+        options.set(
+                
CoreOptions.VARIANT_SHREDDING_ADAPTIVE_MAX_INFER_BUFFER_ROW.key(),
+                String.valueOf(adaptiveSampleRows));
+        options.set(
+                
CoreOptions.VARIANT_SHREDDING_MIN_FIELD_CARDINALITY_RATIO.key(),
+                String.valueOf(admissionRatio));
+        options.set(
+                CoreOptions.VARIANT_SHREDDING_ADAPTIVE_RETENTION_RATIO.key(),
+                String.valueOf(retentionRatio));
+        return options;
+    }
+
+    private Path newFile() {
+        return new Path(parent, UUID.randomUUID() + ".parquet");
+    }
+
     protected List<InternalRow> readRows(ParquetFileFormat format, RowType 
rowType)
             throws IOException {
+        return readRows(format, rowType, file);
+    }
+
+    private List<InternalRow> readRows(ParquetFileFormat format, RowType 
rowType, Path dataFile)
+            throws IOException {
         List<InternalRow> result = new ArrayList<>();
         try (RecordReader<InternalRow> reader =
                 format.createReaderFactory(rowType, rowType, new ArrayList<>())
                         .createReader(
-                                new FormatReaderContext(fileIO, file, 
fileIO.getFileSize(file)))) {
+                                new FormatReaderContext(
+                                        fileIO, dataFile, 
fileIO.getFileSize(dataFile)))) {
             InternalRowSerializer serializer = new 
InternalRowSerializer(rowType);
             reader.forEachRemaining(row -> result.add(serializer.copy(row)));
         }
         return result;
     }
 
+    private List<String> readVariantJson(
+            ParquetFileFormat format, RowType rowType, Path dataFile, int 
fieldIndex)
+            throws IOException {
+        List<String> result = new ArrayList<>();
+        for (InternalRow row : readRows(format, rowType, dataFile)) {
+            result.add(row.getVariant(fieldIndex).toJson());
+        }
+        return result;
+    }
+
     protected void writeRows(FormatWriterFactory factory, InternalRow... rows) 
throws IOException {
+        writeRows(factory, file, rows);
+    }
+
+    private void writeRows(FormatWriterFactory factory, Path dataFile, 
InternalRow... rows)
+            throws IOException {
         FormatWriter writer;
         PositionOutputStream out = null;
         if (factory instanceof SupportsDirectWrite) {
-            writer = ((SupportsDirectWrite) factory).create(fileIO, file, 
"zstd");
+            writer = ((SupportsDirectWrite) factory).create(fileIO, dataFile, 
"zstd");
         } else {
-            out = fileIO.newOutputStream(file, false);
+            out = fileIO.newOutputStream(dataFile, false);
             writer = factory.create(out, "zstd");
         }
         for (InternalRow row : rows) {
@@ -544,6 +828,19 @@ public class InferVariantShreddingWriteTest {
         }
     }
 
+    private RowType readVariantFileType(Path dataFile, String... fieldPath) 
throws IOException {
+        try (ParquetFileReader reader =
+                ParquetUtil.getParquetReader(
+                        fileIO, dataFile, fileIO.getFileSize(dataFile), new 
Options())) {
+            Type variantType = 
reader.getFooter().getFileMetaData().getSchema();
+            for (String fieldName : fieldPath) {
+                variantType = variantType.asGroupType().getType(fieldName);
+            }
+            return VariantMetadataUtils.addVariantMetadata(
+                    
VariantShreddingReadPlanFactory.variantFileType(variantType));
+        }
+    }
+
     protected void verifyShreddingSchema(RowType... expectShreddedTypes) 
throws IOException {
         try (ParquetFileReader reader =
                 ParquetUtil.getParquetReader(

Reply via email to