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 03d7ffd8f4d Send the table config and schema refresh message when a 
schema is overridden (#19677)
03d7ffd8f4d is described below

commit 03d7ffd8f4d9def62143009d421526d141de1f98
Author: Xiaotian (Jackie) Jiang <[email protected]>
AuthorDate: Fri Sep 25 17:52:04 2026 -0700

    Send the table config and schema refresh message when a schema is 
overridden (#19677)
---
 .../helix/core/PinotHelixResourceManager.java      | 50 ++++++++++++++--------
 .../PinotHelixResourceManagerStatelessTest.java    | 32 ++++++++++++++
 2 files changed, 64 insertions(+), 18 deletions(-)

diff --git 
a/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManager.java
 
b/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManager.java
index 89a6f72cea6..62d25ed4440 100644
--- 
a/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManager.java
+++ 
b/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManager.java
@@ -1678,6 +1678,7 @@ public class PinotHelixResourceManager {
       // Update existing schema
       if (override) {
         updateSchema(schema, oldSchema, force);
+        refreshTablesUsingSchema(schemaName);
       } else {
         throw new SchemaAlreadyExistsException("Schema: " + schemaName + " 
already exists");
       }
@@ -1710,30 +1711,43 @@ public class PinotHelixResourceManager {
     }
 
     updateSchema(schema, oldSchema, forceTableSchemaUpdate);
+    if (reload) {
+      reloadTablesUsingSchema(schemaName);
+    } else {
+      refreshTablesUsingSchema(schemaName);
+    }
+  }
+
+  /// Reloads all segments of the tables using the given schema. Logical table 
schemas are skipped because no
+  /// segments are backed by them.
+  private void reloadTablesUsingSchema(String schemaName)
+      throws TableNotFoundException {
     if (ZKMetadataProvider.isLogicalTableExists(_propertyStore, schemaName)) {
-      // For logical table schemas, we do not need to reload segments or send 
schema refresh messages
-      LOGGER.info("Logical table schema: {} updated, no need to reload 
segments or send schema refresh messages",
-          schemaName);
+      LOGGER.info("Logical table schema: {} updated, no need to reload 
segments", schemaName);
       return;
     }
+    LOGGER.info("Reloading tables with name: {}", schemaName);
+    for (String tableNameWithType : getExistingTableNamesWithType(schemaName, 
null)) {
+      reloadAllSegments(tableNameWithType, false, null);
+    }
+  }
+
+  /// Sends the table config and schema refresh message to the servers hosting 
the tables using the given schema.
+  /// Servers reuse their cached schema for ordinary segment loads and only 
refresh it on this message, an explicit
+  /// reload or a table config update, so every schema update has to send it 
or newly loaded segments keep being
+  /// processed against the previous schema. Logical table schemas are skipped 
because no segments are backed by them.
+  private void refreshTablesUsingSchema(String schemaName) {
+    if (ZKMetadataProvider.isLogicalTableExists(_propertyStore, schemaName)) {
+      LOGGER.info("Logical table schema: {} updated, no need to send schema 
refresh messages", schemaName);
+      return;
+    }
+    LOGGER.info("Refreshing schema for tables with name: {}", schemaName);
     try {
-      List<String> tableNamesWithType = 
getExistingTableNamesWithType(schemaName, null);
-      if (reload) {
-        LOGGER.info("Reloading tables with name: {}", schemaName);
-        for (String tableNameWithType : tableNamesWithType) {
-          reloadAllSegments(tableNameWithType, false, null);
-        }
-      } else {
-        LOGGER.info("Refreshing schema for tables with name: {}", schemaName);
-        for (String tableNameWithType : tableNamesWithType) {
-          sendTableConfigSchemaRefreshMessage(tableNameWithType);
-        }
+      for (String tableNameWithType : 
getExistingTableNamesWithType(schemaName, null)) {
+        sendTableConfigSchemaRefreshMessage(tableNameWithType);
       }
     } catch (TableNotFoundException e) {
-      if (reload) {
-        throw e;
-      }
-      // We don't throw exception if no tables found for schema when reload is 
false. Since this could be valid case
+      // A schema can exist before any table uses it
       LOGGER.warn("No tables found for schema (refresh only): {}", schemaName, 
e);
     }
   }
diff --git 
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManagerStatelessTest.java
 
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManagerStatelessTest.java
index 64a8ad22274..d2d2dbc757f 100644
--- 
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManagerStatelessTest.java
+++ 
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/PinotHelixResourceManagerStatelessTest.java
@@ -77,8 +77,10 @@ import 
org.apache.pinot.spi.config.table.ingestion.IngestionConfig;
 import org.apache.pinot.spi.config.tenant.Tenant;
 import org.apache.pinot.spi.config.tenant.TenantRole;
 import org.apache.pinot.spi.data.DateTimeFieldSpec;
+import org.apache.pinot.spi.data.DimensionFieldSpec;
 import org.apache.pinot.spi.data.FieldSpec;
 import org.apache.pinot.spi.data.LogicalTableConfig;
+import org.apache.pinot.spi.data.Schema;
 import org.apache.pinot.spi.stream.LongMsgOffset;
 import org.apache.pinot.spi.stream.PartitionGroupConsumptionStatus;
 import org.apache.pinot.spi.stream.PartitionGroupMetadata;
@@ -1131,6 +1133,36 @@ public class PinotHelixResourceManagerStatelessTest 
extends ControllerTest {
     }
   }
 
+  @Test
+  public void testAddSchemaWithOverrideRefreshesServerCaches()
+      throws Exception {
+    String rawTableName = "schemaOverrideRefreshTest";
+    String offlineTableName = 
TableNameBuilder.OFFLINE.tableNameWithType(rawTableName);
+    addDummySchema(rawTableName);
+    TableConfig tableConfig = new TableConfigBuilder(TableType.OFFLINE)
+        .setTableName(rawTableName)
+        .setBrokerTenant(BROKER_TENANT_NAME)
+        .setServerTenant(SERVER_TENANT_NAME)
+        .build();
+    waitForEVToDisappear(tableConfig.getTableName());
+    _helixResourceManager.addTable(tableConfig);
+
+    PinotHelixResourceManager resourceManager = spy(_helixResourceManager);
+    
doNothing().when(resourceManager).sendTableConfigSchemaRefreshMessage(offlineTableName);
+
+    try {
+      Schema schema = createDummySchema(rawTableName);
+      schema.addField(new DimensionFieldSpec("dimC", 
FieldSpec.DataType.STRING, true));
+      resourceManager.addSchema(schema, true, false);
+
+      assertEquals(resourceManager.getSchema(rawTableName), schema);
+      
verify(resourceManager).sendTableConfigSchemaRefreshMessage(offlineTableName);
+    } finally {
+      _helixResourceManager.deleteOfflineTable(rawTableName);
+      deleteSchema(rawTableName);
+    }
+  }
+
   /// Tests the code path where a subset of merged segments (from the original 
segmentsTo list)
   /// is passed to the endReplace API.
   /// @throws Exception


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

Reply via email to