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]