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())))