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

Jackie-Jiang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git


The following commit(s) were added to refs/heads/master by this push:
     new 6e38e4f8c22 Disallow numPartitions change on segmentPartitionConfig 
for upsert/dedup (#18868)
6e38e4f8c22 is described below

commit 6e38e4f8c222c377669c3b091147722dfb82e120
Author: Chaitanya Deepthi <[email protected]>
AuthorDate: Mon Aug 10 10:15:07 2026 -0700

    Disallow numPartitions change on segmentPartitionConfig for upsert/dedup 
(#18868)
---
 .../segment/local/utils/TableConfigUtils.java      | 33 ++++++++++++++
 .../segment/local/utils/TableConfigUtilsTest.java  | 50 ++++++++++++++++++++++
 2 files changed, 83 insertions(+)

diff --git 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/utils/TableConfigUtils.java
 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/utils/TableConfigUtils.java
index 222a8396e78..fc9158af589 100644
--- 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/utils/TableConfigUtils.java
+++ 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/utils/TableConfigUtils.java
@@ -1474,6 +1474,7 @@ public final class TableConfigUtils {
     List<String> violations = new ArrayList<>();
     validateUpsertConfigUpdate(newConfig, existingConfig, violations);
     validateDedupConfigUpdate(newConfig, existingConfig, violations);
+    validatePartitionConfigUpdate(newConfig, existingConfig, violations);
     validateMaterializedViewConfigUpdate(newConfig, existingConfig, 
violations);
 
     return violations;
@@ -1628,6 +1629,38 @@ public final class TableConfigUtils {
     }
   }
 
+  /// On upsert/dedup tables, rejects `numPartitions` changes on 
already-partitioned columns — the one
+  /// segmentPartitionConfig change that silently breaks broker query pruning 
against existing segments.
+  /// Adding, removing, and `functionName` / `functionConfig` changes are all 
allowed. Bypassable with force update.
+  private static void validatePartitionConfigUpdate(TableConfig newConfig, 
TableConfig existingConfig,
+      List<String> violations) {
+    if (!(existingConfig.isUpsertEnabled() || existingConfig.isDedupEnabled())
+        || !(newConfig.isUpsertEnabled() || newConfig.isDedupEnabled())) {
+      return;
+    }
+    Map<String, ColumnPartitionConfig> existingMap = 
getColumnPartitionMap(existingConfig);
+    Map<String, ColumnPartitionConfig> newMap = 
getColumnPartitionMap(newConfig);
+    for (Map.Entry<String, ColumnPartitionConfig> entry : 
existingMap.entrySet()) {
+      ColumnPartitionConfig newCol = newMap.get(entry.getKey());
+      if (newCol != null && entry.getValue().getNumPartitions() != 
newCol.getNumPartitions()) {
+        violations.add(String.format(
+            "segmentPartitionConfig numPartitions cannot change for 
upsert/dedup column '%s' (%d -> %d)",
+            entry.getKey(), entry.getValue().getNumPartitions(), 
newCol.getNumPartitions()));
+      }
+    }
+  }
+
+  private static Map<String, ColumnPartitionConfig> 
getColumnPartitionMap(TableConfig tableConfig) {
+    if (tableConfig.getIndexingConfig() == null) {
+      return Map.of();
+    }
+    SegmentPartitionConfig partitionConfig = 
tableConfig.getIndexingConfig().getSegmentPartitionConfig();
+    if (partitionConfig == null || partitionConfig.getColumnPartitionMap() == 
null) {
+      return Map.of();
+    }
+    return partitionConfig.getColumnPartitionMap();
+  }
+
   private static void validateMaterializedViewConfigUpdate(TableConfig 
newConfig, TableConfig existingConfig,
       List<String> violations) {
     if (existingConfig.isMaterializedView() != newConfig.isMaterializedView()) 
{
diff --git 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/utils/TableConfigUtilsTest.java
 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/utils/TableConfigUtilsTest.java
index 2354cd03486..310bbfadbc2 100644
--- 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/utils/TableConfigUtilsTest.java
+++ 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/utils/TableConfigUtilsTest.java
@@ -4528,6 +4528,56 @@ public class TableConfigUtilsTest {
     assertTrue(violations.isEmpty(), "Expected no violations for non-upsert 
tables, but got: " + violations);
   }
 
+  private static TableConfig upsertTable(@Nullable SegmentPartitionConfig 
partition) {
+    return new 
TableConfigBuilder(TableType.REALTIME).setTableName(TABLE_NAME).setTimeColumnName(TIME_COLUMN)
+        .setUpsertConfig(new 
UpsertConfig(UpsertConfig.Mode.FULL)).setSegmentPartitionConfig(partition).build();
+  }
+
+  private static SegmentPartitionConfig partition(String col, String fn, int 
n) {
+    return new SegmentPartitionConfig(Map.of(col, new 
ColumnPartitionConfig(fn, n)));
+  }
+
+  @Test
+  public void testNumPartitionsChangeRejected() {
+    List<String> violations = TableConfigUtils.validateBackwardCompatibility(
+        upsertTable(partition("myCol", "Murmur", 8)), 
upsertTable(partition("myCol", "Murmur", 4)));
+    assertEquals(violations.size(), 1);
+    assertTrue(violations.get(0).contains("numPartitions"));
+  }
+
+  @Test
+  public void testDedupNumPartitionsChangeRejected() {
+    TableConfig existing = new 
TableConfigBuilder(TableType.REALTIME).setTableName(TABLE_NAME)
+        .setTimeColumnName(TIME_COLUMN).setDedupConfig(new DedupConfig())
+        .setSegmentPartitionConfig(partition("myCol", "Murmur", 4)).build();
+    TableConfig updated = new 
TableConfigBuilder(TableType.REALTIME).setTableName(TABLE_NAME)
+        .setTimeColumnName(TIME_COLUMN).setDedupConfig(new DedupConfig())
+        .setSegmentPartitionConfig(partition("myCol", "Murmur", 8)).build();
+    assertEquals(TableConfigUtils.validateBackwardCompatibility(updated, 
existing).size(), 1);
+  }
+
+  @Test
+  public void testAddAndRemovePartitionAllowed() {
+    assertTrue(TableConfigUtils.validateBackwardCompatibility(
+        upsertTable(partition("myCol", "Murmur", 4)), 
upsertTable(null)).isEmpty());
+    assertTrue(TableConfigUtils.validateBackwardCompatibility(
+        upsertTable(null), upsertTable(partition("myCol", "Murmur", 
4))).isEmpty());
+  }
+
+  @Test
+  public void testFunctionNameChangeAllowed() {
+    assertTrue(TableConfigUtils.validateBackwardCompatibility(
+        upsertTable(partition("myCol", "Modulo", 4)), 
upsertTable(partition("myCol", "Murmur", 4))).isEmpty());
+  }
+
+  @Test
+  public void testNonUpsertUnchecked() {
+    TableConfig plain = new 
TableConfigBuilder(TableType.REALTIME).setTableName(TABLE_NAME)
+        
.setTimeColumnName(TIME_COLUMN).setSegmentPartitionConfig(partition("myCol", 
"Murmur", 8)).build();
+    assertTrue(TableConfigUtils.validateBackwardCompatibility(plain, 
upsertTable(partition("myCol", "Murmur", 4)))
+        .isEmpty());
+  }
+
   @Test
   public void 
testValidateMaterializedViewInvariantsFlagRequiresOfflineAndTask() {
     TableConfig mvWithoutTask = new 
TableConfigBuilder(TableType.OFFLINE).setTableName("mv_test")


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to