From 2c1db90fc1f726f1952c48fd920bfd06814f2109 Mon Sep 17 00:00:00 2001 From: rahulsmahadev Date: Thu, 6 Aug 2026 01:34:53 +0000 Subject: [PATCH] Core: Add max-file-group-input-files to valid rewrite options max-file-group-input-files is read and validated by SizeBasedFileRewritePlanner but was missing from validOptions(), so RewriteDataFilesSparkAction rejected it with "Cannot use options [max-file-group-input-files]". Flink exposes the option through RewriteDataFiles.maxFileGroupInputFiles() where it works, while Spark rejected it. The option was added in #14837 without updating validOptions(). Add the option to validOptions() and cover it with a rewrite that sets it, which fails without this change, plus a planner test that it caps the number of input files per group. --- .../actions/SizeBasedFileRewritePlanner.java | 3 +- .../TestBinPackRewriteFilePlanner.java | 36 +++++++++++++++++++ ...stBinPackRewritePositionDeletePlanner.java | 1 + .../TestSizeBasedFileRewritePlanner.java | 3 +- .../actions/TestRewriteDataFilesAction.java | 20 +++++++++++ 5 files changed, 61 insertions(+), 2 deletions(-) diff --git a/core/src/main/java/org/apache/iceberg/actions/SizeBasedFileRewritePlanner.java b/core/src/main/java/org/apache/iceberg/actions/SizeBasedFileRewritePlanner.java index bcd00541308f..8784f41a1c3e 100644 --- a/core/src/main/java/org/apache/iceberg/actions/SizeBasedFileRewritePlanner.java +++ b/core/src/main/java/org/apache/iceberg/actions/SizeBasedFileRewritePlanner.java @@ -149,7 +149,8 @@ public Set validOptions() { MAX_FILE_SIZE_BYTES, MIN_INPUT_FILES, REWRITE_ALL, - MAX_FILE_GROUP_SIZE_BYTES); + MAX_FILE_GROUP_SIZE_BYTES, + MAX_FILE_GROUP_INPUT_FILES); } @Override diff --git a/core/src/test/java/org/apache/iceberg/actions/TestBinPackRewriteFilePlanner.java b/core/src/test/java/org/apache/iceberg/actions/TestBinPackRewriteFilePlanner.java index aa65140c0b89..f7a663f06612 100644 --- a/core/src/test/java/org/apache/iceberg/actions/TestBinPackRewriteFilePlanner.java +++ b/core/src/test/java/org/apache/iceberg/actions/TestBinPackRewriteFilePlanner.java @@ -175,6 +175,41 @@ void testMaxGroupSize() { assertThat(plan.groupsInPartition(FILE_6.partition())).isEqualTo(1); } + @Test + void testMaxFileGroupInputFiles() { + addFiles(); + // First, establish baseline without the constraint + BinPackRewriteFilePlanner baselinePlanner = new BinPackRewriteFilePlanner(table); + baselinePlanner.init(REWRITE_ALL); + FileRewritePlan baselinePlan = + baselinePlanner.plan(); + int baselineGroupCount = baselinePlan.totalGroupCount(); + int baselineGroupsInPartition0 = baselinePlan.groupsInPartition(FILE_1.partition()); + + // Now test with max-file-group-input-files set to 2, which limits each group to 2 files + // Partition 0 has 3 files (FILE_1, FILE_2, FILE_3), so it should be split into 2 groups + // Partition 1 has 2 files (FILE_4, FILE_5), so it should be 1 group + // Partition 2 has 1 file (FILE_6), so it should be 1 group + BinPackRewriteFilePlanner constrainedPlanner = new BinPackRewriteFilePlanner(table); + constrainedPlanner.init( + ImmutableMap.of( + BinPackRewriteFilePlanner.REWRITE_ALL, + "true", + BinPackRewriteFilePlanner.MAX_FILE_GROUP_INPUT_FILES, + "2")); + + FileRewritePlan constrainedPlan = + constrainedPlanner.plan(); + + // Verify the constraint is honored: should have MORE groups when input files are limited + assertThat(constrainedPlan.totalGroupCount()).isGreaterThan(baselineGroupCount).isEqualTo(4); + assertThat(constrainedPlan.groupsInPartition(FILE_1.partition())) + .isGreaterThan(baselineGroupsInPartition0) + .isEqualTo(2); + assertThat(constrainedPlan.groupsInPartition(FILE_4.partition())).isEqualTo(1); + assertThat(constrainedPlan.groupsInPartition(FILE_6.partition())).isEqualTo(1); + } + @Test void testFilter() { addFiles(); @@ -289,6 +324,7 @@ void testValidOptions() { BinPackRewriteFilePlanner.MIN_INPUT_FILES, BinPackRewriteFilePlanner.REWRITE_ALL, BinPackRewriteFilePlanner.MAX_FILE_GROUP_SIZE_BYTES, + BinPackRewriteFilePlanner.MAX_FILE_GROUP_INPUT_FILES, BinPackRewriteFilePlanner.DELETE_FILE_THRESHOLD, BinPackRewriteFilePlanner.DELETE_RATIO_THRESHOLD, RewriteDataFiles.REWRITE_JOB_ORDER, diff --git a/core/src/test/java/org/apache/iceberg/actions/TestBinPackRewritePositionDeletePlanner.java b/core/src/test/java/org/apache/iceberg/actions/TestBinPackRewritePositionDeletePlanner.java index c527b658de0f..d201f880f5e3 100644 --- a/core/src/test/java/org/apache/iceberg/actions/TestBinPackRewritePositionDeletePlanner.java +++ b/core/src/test/java/org/apache/iceberg/actions/TestBinPackRewritePositionDeletePlanner.java @@ -239,6 +239,7 @@ void testValidOptions() { BinPackRewritePositionDeletePlanner.MIN_INPUT_FILES, BinPackRewritePositionDeletePlanner.REWRITE_ALL, BinPackRewritePositionDeletePlanner.MAX_FILE_GROUP_SIZE_BYTES, + BinPackRewritePositionDeletePlanner.MAX_FILE_GROUP_INPUT_FILES, RewriteDataFiles.REWRITE_JOB_ORDER)); } diff --git a/core/src/test/java/org/apache/iceberg/actions/TestSizeBasedFileRewritePlanner.java b/core/src/test/java/org/apache/iceberg/actions/TestSizeBasedFileRewritePlanner.java index 7d24e923387b..84d54cd98b36 100644 --- a/core/src/test/java/org/apache/iceberg/actions/TestSizeBasedFileRewritePlanner.java +++ b/core/src/test/java/org/apache/iceberg/actions/TestSizeBasedFileRewritePlanner.java @@ -103,7 +103,8 @@ void testValidOptions() { BinPackRewriteFilePlanner.MAX_FILE_SIZE_BYTES, BinPackRewriteFilePlanner.MIN_INPUT_FILES, BinPackRewriteFilePlanner.REWRITE_ALL, - BinPackRewriteFilePlanner.MAX_FILE_GROUP_SIZE_BYTES)); + BinPackRewriteFilePlanner.MAX_FILE_GROUP_SIZE_BYTES, + BinPackRewriteFilePlanner.MAX_FILE_GROUP_INPUT_FILES)); } @Test diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java index c0899c9ac77d..032025b49572 100644 --- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java +++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java @@ -1486,6 +1486,26 @@ public void testInvalidOptions() { .hasMessageContaining("requires enabling Iceberg Spark session extensions"); } + @TestTemplate + public void testMaxFileGroupInputFilesOption() { + Table table = createTable(4); + shouldHaveFiles(table, 4); + + List originalData = currentData(); + long dataSizeBefore = testDataSize(table); + + RewriteDataFiles.Result result = + basicRewrite(table) + .option(SizeBasedFileRewritePlanner.MAX_FILE_GROUP_INPUT_FILES, "2") + .execute(); + + assertThat(result.rewriteResults()).as("Action should rewrite file groups").isNotEmpty(); + assertThat(result.rewrittenBytesCount()).isEqualTo(dataSizeBefore); + + table.refresh(); + assertEquals("Rows must match", originalData, currentData()); + } + @TestTemplate public void testSortMultipleGroups() { Table table = createTable(20);