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

danny0405 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new c74939e3b1ed perf(flink): validate schema and table properties once 
per instant (#20122)
c74939e3b1ed is described below

commit c74939e3b1ed026d93426932e82192d9291c4570
Author: Danny Chan <[email protected]>
AuthorDate: Mon Sep 28 18:34:26 2026 +0800

    perf(flink): validate schema and table properties once per instant (#20122)
    
    * perf(flink): validate schema and table properties once per instant
---
 .../org/apache/hudi/cli/commands/SparkMain.java    |   2 +-
 .../apache/hudi/client/BaseHoodieWriteClient.java  |  34 +++----
 .../hudi/client/TestBaseHoodieWriteClient.java     |  58 +++++++----
 .../apache/hudi/client/HoodieFlinkWriteClient.java |  18 ++--
 .../apache/hudi/client/TestFlinkWriteClient.java   |  69 +++++++++++++
 .../hudi/sink/StreamWriteOperatorCoordinator.java  |   6 +-
 .../sink/TestStreamWriteOperatorCoordinator.java   | 113 +++++++++++++++++++++
 .../hudi/utilities/HoodieDropPartitionsTool.java   |   3 +-
 8 files changed, 256 insertions(+), 47 deletions(-)

diff --git a/hudi-cli/src/main/java/org/apache/hudi/cli/commands/SparkMain.java 
b/hudi-cli/src/main/java/org/apache/hudi/cli/commands/SparkMain.java
index 8a64538eadc2..4a0e17605583 100644
--- a/hudi-cli/src/main/java/org/apache/hudi/cli/commands/SparkMain.java
+++ b/hudi-cli/src/main/java/org/apache/hudi/cli/commands/SparkMain.java
@@ -276,7 +276,7 @@ public class SparkMain {
       HoodieWriteConfig config = client.getConfig();
       HoodieEngineContext context = client.getEngineContext();
       HoodieSparkTable table = HoodieSparkTable.create(config, context);
-      
client.validateAgainstTableProperties(table.getMetaClient().getTableConfig(), 
config);
+      client.validateAgainstTableProperties(table.getMetaClient(), config, 
WriteOperationType.UNKNOWN);
       WriteMarkersFactory.get(config.getMarkersType(), table, instantTime)
           .quietDeleteMarkerDir(context, config.getMarkersDeleteParallelism());
       return 0;
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java
index 44a91474e403..004f8a5dce32 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java
@@ -1451,10 +1451,9 @@ public abstract class BaseHoodieWriteClient<T, I, K, O> 
extends BaseHoodieClient
   }
 
   /**
-   * Performs necessary bootstrapping operations (for ex, validating whether 
Metadata Table has to be bootstrapped).
+   * Performs necessary bootstrapping operations and validates the table 
properties.
    *
-   * <p>NOTE: THIS OPERATION IS EXECUTED UNDER LOCK, THEREFORE SHOULD AVOID 
ANY OPERATIONS
-   *          NOT REQUIRING EXTERNAL SYNCHRONIZATION
+   * <p>Upgrade and metadata table initialization execute under lock. Table 
properties are validated afterward.
    *
    * @param metaClient instance of {@link HoodieTableMetaClient}
    * @param instantTime current inflight instant time
@@ -1466,14 +1465,14 @@ public abstract class BaseHoodieWriteClient<T, I, K, O> 
extends BaseHoodieClient
           metaClient.getTableType()), instantTime.get()));
     }
     boolean requiresInitTable = needsUpgrade(metaClient) || 
config.isMetadataTableEnabled();
-    if (!requiresInitTable) {
-      return;
+    if (requiresInitTable) {
+      executeUsingTxnManager(ownerInstant, () -> {
+        tryUpgrade(metaClient, instantTime);
+        // TODO: this also does MT table management..
+        initMetadataTable(instantTime, metaClient);
+      });
     }
-    executeUsingTxnManager(ownerInstant, () -> {
-      tryUpgrade(metaClient, instantTime);
-      // TODO: this also does MT table management..
-      initMetadataTable(instantTime, metaClient);
-    });
+    validateAgainstTableProperties(metaClient, config, operationType);
   }
 
   /**
@@ -1511,14 +1510,8 @@ public abstract class BaseHoodieWriteClient<T, I, K, O> 
extends BaseHoodieClient
     }
 
     doInitTable(operationType, metaClient, instantTime);
-    if (WriteOperationType.isInsert(operationType) || 
WriteOperationType.isChangingRecords(operationType)) {
-      ensureComplexKeyGenEncodingRecorded(metaClient);
-    }
     HoodieTable table = createTable(config, metaClient);
 
-    // Validate table properties
-    validateAgainstTableProperties(table.getMetaClient().getTableConfig(), 
config);
-
     switch (operationType) {
       case INSERT:
       case INSERT_PREPPED:
@@ -1558,7 +1551,8 @@ public abstract class BaseHoodieWriteClient<T, I, K, O> 
extends BaseHoodieClient
   }
 
   /**
-   * Pure validation: this method reads both configs and throws, and never 
modifies either.
+   * Validates the write configuration against the table properties. For 
writes that key records,
+   * ensures the complex key generator encoding is recorded; engines may 
record a missing encoding.
    *
    * <p>Nothing reconciles the write config against the table beforehand. That 
is deliberate: the mode
    * is read further down the write path by handles and writer factories, some 
of which hold no table
@@ -1566,7 +1560,11 @@ public abstract class BaseHoodieWriteClient<T, I, K, O> 
extends BaseHoodieClient
    * in. This gate is what makes that true, by refusing writes whose 
meta-field settings do not already
    * agree with the table.
    */
-  public void validateAgainstTableProperties(HoodieTableConfig tableConfig, 
HoodieWriteConfig writeConfig) {
+  public void validateAgainstTableProperties(HoodieTableMetaClient metaClient, 
HoodieWriteConfig writeConfig, WriteOperationType operationType) {
+    if (WriteOperationType.isInsert(operationType) || 
WriteOperationType.isChangingRecords(operationType)) {
+      ensureComplexKeyGenEncodingRecorded(metaClient);
+    }
+    HoodieTableConfig tableConfig = metaClient.getTableConfig();
     // mismatch of table versions.
     CommonClientUtils.validateTableVersion(tableConfig, writeConfig);
 
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/client/TestBaseHoodieWriteClient.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/client/TestBaseHoodieWriteClient.java
index 31a6ae861f72..2c7fac462e68 100644
--- 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/client/TestBaseHoodieWriteClient.java
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/client/TestBaseHoodieWriteClient.java
@@ -61,6 +61,7 @@ import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.CsvSource;
+import org.junit.jupiter.params.provider.EnumSource;
 import org.junit.jupiter.params.provider.MethodSource;
 import org.mockito.InOrder;
 import org.mockito.Mockito;
@@ -105,6 +106,12 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
         mock(BaseHoodieTableServiceClient.class));
   }
 
+  private static HoodieTableMetaClient 
metaClientWithTableConfig(HoodieTableConfig tableConfig) {
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    return metaClient;
+  }
+
   @Test
   void validateAgainstTablePropertiesRejectsMetaFieldsModeMismatch() throws 
IOException {
     initMetaClient();
@@ -122,7 +129,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
         "precondition: both legacy booleans are false, so only the enum 
comparison can catch this");
 
     HoodieException ex = assertThrows(HoodieException.class, () ->
-        
validatorClient(noneWriteConfig).validateAgainstTableProperties(commitTimeOnlyTable,
 noneWriteConfig));
+        
validatorClient(noneWriteConfig).validateAgainstTableProperties(metaClientWithTableConfig(commitTimeOnlyTable),
 noneWriteConfig, WriteOperationType.INSERT));
     
assertTrue(ex.getMessage().contains(HoodieTableConfig.META_FIELDS_MODE.key()),
         "error must name the mode property: " + ex.getMessage());
     assertTrue(ex.getMessage().contains("COMMIT_TIME_ONLY") && 
ex.getMessage().contains("NONE"),
@@ -143,7 +150,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
         "precondition: on its own an unstated writer resolves to the ALL 
default");
 
     validatorClient(unstated)
-        
.validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.ALL), 
unstated);
+        
.validateAgainstTableProperties(metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.ALL)),
 unstated, WriteOperationType.INSERT);
 
     assertEquals(MetaFieldsMode.ALL, unstated.getMetaFieldsMode());
     assertTrue(unstated.populateMetaFields(),
@@ -167,7 +174,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
 
     HoodieException ex = assertThrows(HoodieException.class, () ->
         validatorClient(unstated).validateAgainstTableProperties(
-            tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY), unstated));
+            
metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY)),
 unstated, WriteOperationType.INSERT));
     
assertTrue(ex.getMessage().contains(HoodieTableConfig.META_FIELDS_MODE.key()), 
ex.getMessage());
     assertTrue(ex.getMessage().contains("COMMIT_TIME_ONLY"),
         "the error must name the mode the writer has to state: " + 
ex.getMessage());
@@ -183,7 +190,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
         .build();
 
     validatorClient(restated)
-        
.validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY),
 restated);
+        
.validateAgainstTableProperties(metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY)),
 restated, WriteOperationType.INSERT);
 
     assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, 
restated.getMetaFieldsMode());
     assertFalse(restated.populateMetaFields(),
@@ -208,7 +215,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
 
     HoodieException ex = assertThrows(HoodieException.class, () ->
         validatorClient(stated).validateAgainstTableProperties(
-            tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY), stated));
+            
metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY)),
 stated, WriteOperationType.INSERT));
     
assertTrue(ex.getMessage().contains(HoodieTableConfig.META_FIELDS_MODE.key()), 
ex.getMessage());
     assertTrue(ex.getMessage().contains("COMMIT_TIME_ONLY") && 
ex.getMessage().contains("NONE"),
         "error must name both modes: " + ex.getMessage());
@@ -228,7 +235,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
 
     assertThrows(HoodieException.class, () ->
         validatorClient(stated).validateAgainstTableProperties(
-            tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME), 
stated));
+            
metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME)),
 stated, WriteOperationType.INSERT));
   }
 
   @Test
@@ -243,9 +250,9 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
         
.withPath(basePath).withMetaFieldsMode(MetaFieldsMode.FILE_NAME_ONLY).build();
 
     assertThrows(HoodieException.class, () -> validatorClient(commitTimeOnly)
-        
.validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.FILE_NAME_ONLY),
 commitTimeOnly));
+        
.validateAgainstTableProperties(metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.FILE_NAME_ONLY)),
 commitTimeOnly, WriteOperationType.INSERT));
     assertThrows(HoodieException.class, () -> validatorClient(fileNameOnly)
-        
.validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY),
 fileNameOnly));
+        
.validateAgainstTableProperties(metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY)),
 fileNameOnly, WriteOperationType.INSERT));
   }
 
   @Test
@@ -260,13 +267,13 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
 
     // Rejected in this direction only because the writer stated the mode 
explicitly.
     assertThrows(HoodieException.class, () -> validatorClient(selective)
-        
.validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.ALL), 
selective));
+        
.validateAgainstTableProperties(metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.ALL)),
 selective, WriteOperationType.INSERT));
     // ...but never as a *widening*: the reverse direction is what isWiderThan 
must catch.
     HoodieWriteConfig allWriter = HoodieWriteConfig.newBuilder()
         .withPath(basePath).withMetaFieldsMode(MetaFieldsMode.ALL).build();
     HoodieException ex = assertThrows(HoodieException.class, () -> 
validatorClient(allWriter)
         .validateAgainstTableProperties(
-            tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME), 
allWriter));
+            
metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME)),
 allWriter, WriteOperationType.INSERT));
     assertTrue(ex.getMessage().contains("would leave earlier commits without 
it"),
         "the message must name the widening as the reason: " + 
ex.getMessage());
   }
@@ -284,7 +291,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
 
     HoodieException ex = assertThrows(HoodieException.class, () ->
         validatorClient(defaultWriteConfig).validateAgainstTableProperties(
-            tableConfigWithMode(MetaFieldsMode.NONE), defaultWriteConfig));
+            
metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.NONE)), 
defaultWriteConfig, WriteOperationType.INSERT));
     assertTrue(ex.getMessage().contains("requests ALL"), ex.getMessage());
     assertTrue(ex.getMessage().contains("NONE"), ex.getMessage());
     // The rejected validation must not have mutated the write config.
@@ -304,7 +311,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
 
     HoodieException ex = assertThrows(HoodieException.class, () ->
         validatorClient(allWriter)
-            
.validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.NONE), 
allWriter));
+            
.validateAgainstTableProperties(metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.NONE)),
 allWriter, WriteOperationType.INSERT));
     assertTrue(ex.getMessage().contains("would leave earlier commits without 
it"),
         "the message must name the widening as the reason: " + 
ex.getMessage());
   }
@@ -323,7 +330,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
 
     HoodieException ex = assertThrows(HoodieException.class, () ->
         validatorClient(legacyFalse)
-            
.validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.ALL), 
legacyFalse));
+            
.validateAgainstTableProperties(metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.ALL)),
 legacyFalse, WriteOperationType.INSERT));
     assertTrue(ex.getMessage().contains("hudi-cli"),
         "the message must point at the sanctioned mutation path: " + 
ex.getMessage());
   }
@@ -335,7 +342,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
     HoodieWriteConfig defaultWriteConfig = 
HoodieWriteConfig.newBuilder().withPath(basePath).build();
     assertEquals(MetaFieldsMode.ALL, defaultWriteConfig.getMetaFieldsMode());
     validatorClient(defaultWriteConfig)
-        
.validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.ALL), 
defaultWriteConfig);
+        
.validateAgainstTableProperties(metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.ALL)),
 defaultWriteConfig, WriteOperationType.INSERT);
 
     // And a selective writer against a table recorded with the same mode.
     HoodieWriteConfig selectiveWriteConfig = HoodieWriteConfig.newBuilder()
@@ -343,7 +350,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
         .withMetaFieldsMode(MetaFieldsMode.COMMIT_TIME_ONLY)
         .build();
     validatorClient(selectiveWriteConfig).validateAgainstTableProperties(
-        tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY), 
selectiveWriteConfig);
+        
metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY)),
 selectiveWriteConfig, WriteOperationType.INSERT);
   }
 
   @Test
@@ -520,6 +527,23 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
     assertEquals(requestedTime, 
writeTimeline.lastInstant().get().requestedTime());
   }
 
+  @ParameterizedTest
+  @EnumSource(value = WriteOperationType.class, names = {"INSERT", 
"INSERT_PREPPED", "UPSERT", "UPSERT_PREPPED",
+      "BULK_INSERT", "BULK_INSERT_PREPPED", "INSERT_OVERWRITE", 
"INSERT_OVERWRITE_TABLE", "DELETE", "DELETE_PREPPED"})
+  void 
testValidateAgainstTablePropertiesRequiresRecordedComplexKeygenEncoding(WriteOperationType
 operationType) throws IOException {
+    initMetaClient();
+    HoodieTableConfig tableConfig = tableConfigWithMode(MetaFieldsMode.ALL);
+    tableConfig.setValue(HoodieTableConfig.KEY_GENERATOR_CLASS_NAME, 
ComplexAvroKeyGenerator.class.getName());
+    tableConfig.setValue(HoodieTableConfig.RECORDKEY_FIELDS, "id");
+    HoodieTableMetaClient tableMetaClient = 
metaClientWithTableConfig(tableConfig);
+    HoodieWriteConfig writeConfig = 
HoodieWriteConfig.newBuilder().withPath(basePath).build();
+    try (BaseHoodieWriteClient<?, ?, ?, ?> writeClient = 
validatorClient(writeConfig)) {
+      HoodieException e = assertThrows(HoodieException.class,
+          () -> writeClient.validateAgainstTableProperties(tableMetaClient, 
writeConfig, operationType));
+      
assertTrue(e.getMessage().contains(HoodieTableConfig.COMPLEX_KEYGEN_ENCODING.key()),
 e.getMessage());
+    }
+  }
+
   /** A write that keys records on a tracked table without the property is 
refused; table services, partition deletes and rollbacks are not. */
   @Test
   void testInitTableRequiresRecordedComplexKeygenEncoding() throws IOException 
{
@@ -743,7 +767,7 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
 
     HoodieException ex = assertThrows(HoodieException.class, () ->
         validatorClient(writeConfig).validateAgainstTableProperties(
-            tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY), 
writeConfig));
+            
metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY)),
 writeConfig, WriteOperationType.INSERT));
     assertTrue(ex.getMessage().contains(indexTypeName),
         "the message must name the index type: " + ex.getMessage());
     
assertTrue(ex.getMessage().contains(HoodieRecord.RECORD_KEY_METADATA_FIELD),
@@ -767,6 +791,6 @@ class TestBaseHoodieWriteClient extends 
HoodieCommonTestHarness {
         .build();
 
     validatorClient(writeConfig).validateAgainstTableProperties(
-        tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY), writeConfig);
+        
metaClientWithTableConfig(tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY)),
 writeConfig, WriteOperationType.INSERT);
   }
 }
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java
index cd677fa34fea..ee4b99dade20 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java
@@ -74,6 +74,9 @@ import java.util.stream.Collectors;
  * <p>The client is used both on driver (for starting/committing transactions)
  * and executor (for writing dataset).
  *
+ * <p>Table property and schema validation for streaming writes is performed 
by the operator coordinator before
+ * making each new instant available to the write tasks.
+ *
  * @param <T> type of the payload
  */
 @Slf4j
@@ -231,7 +234,6 @@ public class HoodieFlinkWriteClient<T>
   public List<WriteStatus> upsert(Iterator<HoodieRecord<T>> records, 
BucketInfo bucketInfo, String instantTime) {
     HoodieTable<T, List<HoodieRecord<T>>, List<HoodieKey>, List<WriteStatus>> 
table =
         initTable(WriteOperationType.UPSERT, Option.ofNullable(instantTime));
-    table.validateUpsertSchema();
 
     preWrite(instantTime, WriteOperationType.UPSERT, table.getMetaClient());
     HoodieWriteMetadata<List<WriteStatus>> result;
@@ -249,7 +251,6 @@ public class HoodieFlinkWriteClient<T>
     // only used for metadata table, the upsert happens in single thread
     HoodieTable<T, List<HoodieRecord<T>>, List<HoodieKey>, List<WriteStatus>> 
table =
         initTable(WriteOperationType.UPSERT, Option.ofNullable(instantTime));
-    table.validateUpsertSchema();
     preWrite(instantTime, WriteOperationType.UPSERT_PREPPED, 
table.getMetaClient(), Option.of(HoodieListData.eager(preppedRecords)));
     Map<String, List<HoodieRecord<T>>> preppedRecordsByFileId = 
preppedRecords.stream().parallel()
         .collect(Collectors.groupingBy(r -> 
r.getCurrentLocation().getFileId()));
@@ -273,7 +274,6 @@ public class HoodieFlinkWriteClient<T>
   public List<WriteStatus> insert(Iterator<HoodieRecord<T>> records, 
BucketInfo bucketInfo, String instantTime) {
     HoodieTable<T, List<HoodieRecord<T>>, List<HoodieKey>, List<WriteStatus>> 
table =
         initTable(WriteOperationType.INSERT, Option.ofNullable(instantTime));
-    table.validateInsertSchema();
 
     preWrite(instantTime, WriteOperationType.INSERT, table.getMetaClient());
     HoodieWriteMetadata<List<WriteStatus>> result;
@@ -290,7 +290,6 @@ public class HoodieFlinkWriteClient<T>
   public List<WriteStatus> insertOverwrite(Iterator<HoodieRecord<T>> records, 
BucketInfo bucketInfo, String instantTime) {
     HoodieTable<T, List<HoodieRecord<T>>, List<HoodieKey>, List<WriteStatus>> 
table =
         initTable(WriteOperationType.INSERT_OVERWRITE, 
Option.ofNullable(instantTime));
-    table.validateInsertSchema();
     preWrite(instantTime, WriteOperationType.INSERT_OVERWRITE, 
table.getMetaClient());
     // create the write handle if not exists
     HoodieWriteMetadata<List<WriteStatus>> result;
@@ -303,7 +302,6 @@ public class HoodieFlinkWriteClient<T>
   @Override
   public List<WriteStatus> insertOverwriteTable(Iterator<HoodieRecord<T>> 
records, BucketInfo bucketInfo, String instantTime) {
     HoodieTable table = initTable(WriteOperationType.INSERT_OVERWRITE_TABLE, 
Option.ofNullable(instantTime));
-    table.validateInsertSchema();
     preWrite(instantTime, WriteOperationType.INSERT_OVERWRITE_TABLE, 
table.getMetaClient());
     // create the write handle if not exists
     HoodieWriteMetadata<List<WriteStatus>> result;
@@ -333,7 +331,6 @@ public class HoodieFlinkWriteClient<T>
     // only used for metadata table, the bulk_insert happens in single JVM 
process
     HoodieTable<T, List<HoodieRecord<T>>, List<HoodieKey>, List<WriteStatus>> 
table =
         initTable(WriteOperationType.BULK_INSERT_PREPPED, 
Option.ofNullable(instantTime));
-    table.validateInsertSchema();
     preWrite(instantTime, WriteOperationType.BULK_INSERT_PREPPED, 
table.getMetaClient(), Option.of(HoodieListData.eager(preppedRecords)));
     Map<String, List<HoodieRecord<T>>> preppedRecordsByFileId = 
preppedRecords.stream().parallel()
         .collect(Collectors.groupingBy(r -> 
r.getCurrentLocation().getFileId()));
@@ -403,10 +400,15 @@ public class HoodieFlinkWriteClient<T>
   }
 
   /**
-   * Refresh the last transaction metadata,
+   * Validate the table properties and write schema, and refresh the last 
transaction metadata,
    * should be called before the Driver starts a new transaction with a 
reloaded metaclient.
    */
   public void preTxn(WriteOperationType operationType, HoodieTableMetaClient 
metaClient) {
+    // Validate once on the coordinator using its refreshed timeline, before 
writers receive the instant.
+    validateAgainstTableProperties(metaClient, config, operationType);
+    if (!metaClient.isMetadataTable() && 
(WriteOperationType.isChangingRecords(operationType) || 
WriteOperationType.isInsert(operationType))) {
+      createTable(config, metaClient).validateSchema();
+    }
     if (txnManager.isLockRequired() && 
config.needResolveWriteConflict(operationType, metaClient.isMetadataTable(), 
config, metaClient.getTableConfig())) {
       this.lastCompletedTxnAndMetadata = 
TransactionUtils.getLastCompletedTxnInstantAndMetadata(metaClient);
       this.pendingInflightAndRequestedInstants = 
TransactionUtils.getInflightAndRequestedInstants(metaClient);
@@ -489,6 +491,8 @@ public class HoodieFlinkWriteClient<T>
     // no need to execute the upgrade/downgrade on each write in streaming.
 
     // flink performs metadata table bootstrap on the coordinator when it 
starts up.
+
+    // flink validates table properties on the coordinator in preTxn.
   }
 
   /**
diff --git 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java
 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java
index bf41c13426bb..00c33091f2cd 100644
--- 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java
+++ 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java
@@ -24,7 +24,10 @@ import org.apache.hudi.common.config.HoodieMetadataConfig;
 import org.apache.hudi.common.engine.EngineType;
 import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy;
 import org.apache.hudi.common.model.HoodieKey;
+import org.apache.hudi.common.model.MetaFieldsMode;
 import org.apache.hudi.common.model.TableServiceType;
+import org.apache.hudi.common.model.WriteOperationType;
+import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.config.HoodieCleanConfig;
@@ -40,6 +43,7 @@ import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
 import org.junit.jupiter.params.provider.ValueSource;
 
 import java.io.IOException;
@@ -51,6 +55,13 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.clearInvocations;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.verify;
 
 public class TestFlinkWriteClient extends HoodieFlinkClientTestHarness {
 
@@ -66,6 +77,64 @@ public class TestFlinkWriteClient extends 
HoodieFlinkClientTestHarness {
     cleanupResources();
   }
 
+  @Test
+  void testTablePropertyValidationRunsOnlyInPreTxn() {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withPath(metaClient.getBasePath())
+        .withEngineType(EngineType.FLINK)
+        .withMetaFieldsMode(MetaFieldsMode.NONE)
+        .withEmbeddedTimelineServerEnabled(false)
+        .build();
+    writeClient = new HoodieFlinkWriteClient(context, writeConfig);
+
+    // Per-bucket table initialization defers validation to the coordinator.
+    writeClient.initTable(WriteOperationType.UPSERT, Option.empty());
+    HoodieException exception = assertThrows(HoodieException.class,
+        () -> writeClient.preTxn(WriteOperationType.UPSERT, metaClient));
+    
assertTrue(exception.getMessage().contains(HoodieTableConfig.META_FIELDS_MODE.key()));
+  }
+
+  @Test
+  void testPreTxnSkipsSchemaValidationForMetadataTableDeletePrepped() {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withPath(metaClient.getBasePath())
+        .withEngineType(EngineType.FLINK)
+        .withEmbeddedTimelineServerEnabled(false)
+        .build();
+    writeClient = spy(new HoodieFlinkWriteClient(context, writeConfig));
+    HoodieTableMetaClient metadataMetaClient = spy(metaClient);
+    doReturn(true).when(metadataMetaClient).isMetadataTable();
+    writeClient.preTxn(WriteOperationType.DELETE_PREPPED, metadataMetaClient);
+
+    verify(writeClient, never()).createTable(writeConfig, metadataMetaClient);
+  }
+
+  @ParameterizedTest
+  @EnumSource(value = WriteOperationType.class, names = {"UPSERT_PREPPED", 
"BULK_INSERT_PREPPED"})
+  void testPreppedSchemaValidationRunsOnlyInPreTxn(WriteOperationType 
operationType) {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withPath(metaClient.getBasePath())
+        .withEngineType(EngineType.FLINK)
+        .withEmbeddedTimelineServerEnabled(false)
+        .build();
+    writeClient = spy(new HoodieFlinkWriteClient(context, writeConfig));
+    HoodieTable table = spy(writeClient.getHoodieTable(false));
+    doReturn(table).when(writeClient).createTable(eq(writeConfig), 
any(HoodieTableMetaClient.class));
+
+    writeClient.preTxn(operationType, metaClient);
+    verify(writeClient).validateAgainstTableProperties(metaClient, 
writeConfig, operationType);
+    verify(table).validateSchema();
+
+    clearInvocations(table, writeClient);
+    if (operationType == WriteOperationType.UPSERT_PREPPED) {
+      writeClient.upsertPreppedRecords(Collections.emptyList(), "001");
+    } else {
+      writeClient.bulkInsertPreppedRecords(Collections.emptyList(), "001", 
Option.empty());
+    }
+    verify(table, never()).validateSchema();
+    verify(writeClient, never()).validateAgainstTableProperties(any(), any(), 
any());
+  }
+
   @ParameterizedTest
   @ValueSource(booleans = {true, false})
   public void testWriteClientAndTableServiceClientWithTimelineServer(
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
index 11fccc32e87f..3100cb5f3cac 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
@@ -518,9 +518,9 @@ public class StreamWriteOperatorCoordinator
   }
 
   private String startInstant() {
-    // refresh the meta client which is reused
-    metaClient.reloadActiveTimeline();
-    // refresh the last txn metadata
+    // Refresh table properties and index definitions as well as the timeline 
before validating the new write.
+    this.metaClient = HoodieTableMetaClient.reload(this.metaClient);
+    // Validate the write and refresh the last txn metadata.
     this.writeClient.preTxn(tableState.operationType, this.metaClient);
     // put the assignment in front of metadata generation,
     // because the instant request from write task is asynchronous.
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
index bd661a9fdbae..12adab2cc3e0 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
@@ -18,18 +18,25 @@
 
 package org.apache.hudi.sink;
 
+import org.apache.hudi.client.HoodieFlinkWriteClient;
 import org.apache.hudi.client.WriteStatus;
 import org.apache.hudi.client.common.HoodieFlinkEngineContext;
 import org.apache.hudi.client.heartbeat.HoodieHeartbeatClient;
 import org.apache.hudi.common.fs.FSUtils;
+import org.apache.hudi.common.model.HoodieCommitMetadata;
 import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy;
+import org.apache.hudi.common.model.HoodieIndexDefinition;
 import org.apache.hudi.common.model.HoodieTableType;
 import org.apache.hudi.common.model.HoodieWriteStat;
+import org.apache.hudi.common.model.MetaFieldsMode;
 import org.apache.hudi.common.model.WriteConcurrencyMode;
+import org.apache.hudi.common.model.WriteOperationType;
+import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.timeline.HoodieActiveTimeline;
 import org.apache.hudi.common.table.timeline.HoodieInstant;
 import org.apache.hudi.common.table.timeline.HoodieTimeline;
+import org.apache.hudi.common.testutils.HoodieTestTable;
 import org.apache.hudi.common.testutils.HoodieTestUtils;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.SerializationUtils;
@@ -39,7 +46,10 @@ import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.configuration.FlinkOptions;
 import org.apache.hudi.configuration.HadoopConfigurations;
 import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.exception.MissingSchemaFieldException;
+import org.apache.hudi.exception.SchemaCompatibilityException;
 import org.apache.hudi.hadoop.fs.HadoopFSUtils;
+import org.apache.hudi.metadata.HoodieIndexVersion;
 import org.apache.hudi.metadata.HoodieTableMetadata;
 import org.apache.hudi.metadata.MetadataPartitionType;
 import org.apache.hudi.sink.event.Correspondent;
@@ -87,6 +97,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Properties;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.TimeUnit;
@@ -144,6 +155,108 @@ public class TestStreamWriteOperatorCoordinator {
     assertNotEquals(instant, inflight, "Should start a new instant");
   }
 
+  @ParameterizedTest
+  @ValueSource(strings = {"upsert", "insert", "bulk_insert", 
"insert_overwrite", "insert_overwrite_table", "delete"})
+  void testValidationOncePerInstant(String operation) throws Exception {
+    coordinator.close();
+    Configuration conf = 
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
+    conf.set(FlinkOptions.OPERATION, operation);
+    coordinator = startCoordinator(conf, 2);
+    HoodieFlinkWriteClient writeClient = 
Mockito.spy(coordinator.getWriteClient());
+    Field writeClientField = 
StreamWriteOperatorCoordinator.class.getDeclaredField("writeClient");
+    writeClientField.setAccessible(true);
+    writeClientField.set(coordinator, writeClient);
+
+    String firstInstant = requestInstantTime(1);
+    assertEquals(firstInstant, requestInstantTime(1));
+    Mockito.verify(writeClient, 
Mockito.times(1)).preTxn(Mockito.eq(WriteOperationType.fromValue(operation)), 
Mockito.any());
+    Mockito.verify(writeClient, 
Mockito.times(1)).validateAgainstTableProperties(Mockito.any(), Mockito.any(), 
Mockito.eq(WriteOperationType.fromValue(operation)));
+
+    // Complete the first checkpoint so the next request can create a new 
instant.
+    coordinator.handleEventFromOperator(0, createOperatorEvent(0, 1, 
firstInstant, "par1", false, true, 0.1));
+    coordinator.handleEventFromOperator(1, createOperatorEvent(1, 1, 
firstInstant, "par2", false, true, 0.2));
+    coordinator.notifyCheckpointComplete(2);
+    assertNull(coordinator.getEventBuffer(1));
+    String nextInstant = requestInstantTime(2);
+    assertNotEquals(firstInstant, nextInstant);
+    Mockito.verify(writeClient, 
Mockito.times(2)).preTxn(Mockito.eq(WriteOperationType.fromValue(operation)), 
Mockito.any());
+    Mockito.verify(writeClient, 
Mockito.times(2)).validateAgainstTableProperties(Mockito.any(), Mockito.any(), 
Mockito.eq(WriteOperationType.fromValue(operation)));
+    assertFalse(((MockOperatorCoordinatorContext) 
coordinator.getContext()).isJobFailed());
+  }
+
+  @ParameterizedTest
+  @ValueSource(strings = {"upsert", "insert", "bulk_insert", 
"insert_overwrite", "insert_overwrite_table", "delete"})
+  void testColumnDropFailsBeforeInstantIsPublished(String operation) throws 
Exception {
+    coordinator.close();
+    Configuration conf = 
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
+    conf.set(FlinkOptions.OPERATION, operation);
+    coordinator = startCoordinator(conf, 2);
+    HoodieWriteConfig writeConfig = coordinator.getWriteClient().getConfig();
+    assertFalse(writeConfig.shouldValidateAvroSchema());
+    assertFalse(writeConfig.shouldAllowAutoEvolutionColumnDrop());
+
+    // Complete an instant using the original schema, then drop a column for 
the next instant.
+    HoodieCommitMetadata metadata = new HoodieCommitMetadata();
+    metadata.addMetadata(HoodieCommitMetadata.SCHEMA_KEY, 
writeConfig.getSchema());
+    HoodieTestTable.of(StreamerUtil.createMetaClient(conf)).addCommit("001", 
Option.of(metadata));
+    conf.set(FlinkOptions.SOURCE_AVRO_SCHEMA_PATH,
+        
getClass().getClassLoader().getResource("test_read_schema_dropped_age.avsc").toString());
+    writeConfig.setSchema(StreamerUtil.getSourceSchema(conf).toString());
+
+    assertInstantCreationFails(conf, MissingSchemaFieldException.class, "age");
+  }
+
+  @Test
+  void testTablePropertyChangesAreValidatedBeforeInstantIsPublished() throws 
Exception {
+    Configuration conf = 
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
+    HoodieTableMetaClient metaClient = StreamerUtil.createMetaClient(conf);
+    Properties updatedProperties = new Properties();
+    updatedProperties.setProperty(HoodieTableConfig.META_FIELDS_MODE.key(), 
MetaFieldsMode.NONE.name());
+    
updatedProperties.setProperty(HoodieTableConfig.POPULATE_META_FIELDS.key(), 
"false");
+    HoodieTableConfig.update(metaClient.getStorage(), 
metaClient.getMetaPath(), updatedProperties);
+
+    assertInstantCreationFails(conf, HoodieException.class, 
HoodieTableConfig.META_FIELDS_MODE.key());
+  }
+
+  @Test
+  void testNewSecondaryIndexIsValidatedBeforeInstantIsPublished() throws 
Exception {
+    Configuration conf = 
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
+    HoodieWriteConfig writeConfig = coordinator.getWriteClient().getConfig();
+    HoodieTableMetaClient metaClient = StreamerUtil.createMetaClient(conf);
+    HoodieCommitMetadata metadata = new HoodieCommitMetadata();
+    metadata.addMetadata(HoodieCommitMetadata.SCHEMA_KEY, 
writeConfig.getSchema());
+    HoodieTestTable.of(metaClient).addCommit("001", Option.of(metadata));
+    // Simulate an index created after the coordinator has loaded its meta 
client.
+    metaClient.buildIndexDefinition(HoodieIndexDefinition.newBuilder()
+        .withIndexName("secondary_index_age")
+        .withIndexType("secondary_index")
+        .withSourceFields(Collections.singletonList("age"))
+        .withVersion(HoodieIndexVersion.V1)
+        .build());
+    String evolvedSchema = writeConfig.getSchema().replace("\"int\"", 
"\"long\"");
+    assertNotEquals(writeConfig.getSchema(), evolvedSchema);
+    writeConfig.setSchema(evolvedSchema);
+
+    assertInstantCreationFails(conf, SchemaCompatibilityException.class, 
"secondary_index_age");
+  }
+
+  private void assertInstantCreationFails(Configuration conf, Class<? extends 
Throwable> causeType, String message) throws Exception {
+    Correspondent.InstantTimeResponse response = 
CoordinationResponseSerDe.unwrap(
+        
coordinator.handleCoordinationRequest(Correspondent.InstantTimeRequest.getInstance(1))
+            .get(1, TimeUnit.SECONDS));
+    assertNull(response.getInstant());
+    assertNull(coordinator.getEventBuffer(1));
+    MockOperatorCoordinatorContext context = (MockOperatorCoordinatorContext) 
coordinator.getContext();
+    assertTrue(context.isJobFailed());
+    Throwable cause = context.getJobFailureReason();
+    while (cause != null && !(causeType.isInstance(cause) && 
cause.getMessage() != null && cause.getMessage().contains(message))) {
+      cause = cause.getCause();
+    }
+    assertNotNull(cause, "Expected " + causeType.getSimpleName() + " 
containing: " + message);
+    HoodieTimeline pending = 
StreamerUtil.createMetaClient(conf).reloadActiveTimeline().filterPendingExcludingCompaction();
+    assertTrue(pending.empty(), "Validation must fail before a new instant is 
created");
+  }
+
   @Test
   public void testTableInitialized() throws IOException {
     final org.apache.hadoop.conf.Configuration hadoopConf = 
HadoopConfigurations.getHadoopConf(new Configuration());
diff --git 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieDropPartitionsTool.java
 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieDropPartitionsTool.java
index 0999f16ed438..4f37ea81fca3 100644
--- 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieDropPartitionsTool.java
+++ 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieDropPartitionsTool.java
@@ -22,6 +22,7 @@ import org.apache.hudi.client.HoodieWriteResult;
 import org.apache.hudi.client.SparkRDDWriteClient;
 import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.model.HoodieRecordPayload;
+import org.apache.hudi.common.model.WriteOperationType;
 import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.timeline.HoodieTimeline;
@@ -327,7 +328,7 @@ public class HoodieDropPartitionsTool implements 
Serializable {
   public void dryRun() {
     try (SparkRDDWriteClient<HoodieRecordPayload> client =  
UtilHelpers.createHoodieClient(jsc, cfg.basePath, "", cfg.parallelism, 
Option.empty(), props)) {
       HoodieSparkTable<HoodieRecordPayload> table = 
HoodieSparkTable.create(client.getConfig(), client.getEngineContext());
-      
client.validateAgainstTableProperties(table.getMetaClient().getTableConfig(), 
client.getConfig());
+      client.validateAgainstTableProperties(table.getMetaClient(), 
client.getConfig(), WriteOperationType.DELETE_PARTITION);
       List<String> parts = Arrays.asList(cfg.partitions.split(","));
       Map<String, List<String>> partitionToReplaceFileIds = 
jsc.parallelize(parts, parts.size()).distinct()
           .mapToPair(partitionPath -> new Tuple2<>(partitionPath, 
table.getSliceView().getLatestFileSlices(partitionPath).map(fg -> 
fg.getFileId()).distinct().collect(Collectors.toList())))

Reply via email to