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]