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 a3229c5bb72f9..78919746219c1 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 @@ -45,6 +45,7 @@ import org.apache.hudi.common.model.HoodieRecord; 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.TableServiceType; import org.apache.hudi.common.model.WriteOperationType; import org.apache.hudi.common.schema.HoodieSchema; @@ -1548,9 +1549,40 @@ public void validateAgainstTableProperties(HoodieTableConfig tableConfig, Hoodie // mismatch of table versions. CommonClientUtils.validateTableVersion(tableConfig, writeConfig); - // Once meta fields are disabled, it cant be re-enabled for a given table. - if (!tableConfig.populateMetaFields() && writeConfig.populateMetaFields()) { - throw new HoodieException(HoodieTableConfig.POPULATE_META_FIELDS.key() + " already disabled for the table. Can't be re-enabled back"); + // Meta-field population is physical, so a writer must not claim columns the table does not + // have. Compare the full enum rather than the legacy booleans: those collapse every selective + // mode to false, so a writer claiming COMMIT_TIME_ONLY against a NONE table would slip through + // and advertise commit times that were never written. + // + // Two distinct cases, because writers routinely omit meta-field settings entirely: + // + // - Widening is always rejected. Enabling a column now would leave earlier commits without it, + // and readers cannot tell the two apart. + // - Any disagreement is rejected when the writer *explicitly* sets hoodie.meta.fields.mode. + // That covers narrowing too, e.g. an explicit NONE against a COMMIT_TIME_ONLY table, which + // would write null commit times while the table still advertises COMMIT_TIME_ONLY and make + // incremental queries silently miss those rows. + // + // A writer that never mentions the mode is left alone: resolving to NONE against an ALL table + // is long-standing behavior for callers that build a write config without restating the table's + // settings, and writing fewer meta columns cannot make a reader believe in absent data. + MetaFieldsMode tableMetaFieldsMode = tableConfig.getMetaFieldsMode(); + MetaFieldsMode writeMetaFieldsMode = writeConfig.getMetaFieldsMode(); + boolean writerStatedMode = writeConfig.contains(HoodieTableConfig.META_FIELDS_MODE) + && !StringUtils.isNullOrEmpty(writeConfig.getString(HoodieTableConfig.META_FIELDS_MODE)); + if (writeMetaFieldsMode.isWiderThan(tableMetaFieldsMode)) { + throw new HoodieException(String.format( + "%s cannot be widened for an existing table: table is %s but the writer requests %s. Meta " + + "columns are physical, so enabling one now would leave earlier commits without it. " + + "Set %s=%s on the writer, or recreate the table to change it.", + HoodieTableConfig.META_FIELDS_MODE.key(), tableMetaFieldsMode, writeMetaFieldsMode, + HoodieTableConfig.META_FIELDS_MODE.key(), tableMetaFieldsMode)); + } else if (writerStatedMode && writeMetaFieldsMode != tableMetaFieldsMode) { + throw new HoodieException(String.format( + "%s mismatch: table is %s but the writer explicitly requests %s. Meta columns are physical, " + + "so the writer must match the table. Set %s=%s on the writer, or recreate the table.", + HoodieTableConfig.META_FIELDS_MODE.key(), tableMetaFieldsMode, writeMetaFieldsMode, + HoodieTableConfig.META_FIELDS_MODE.key(), tableMetaFieldsMode)); } // Meta fields can be disabled only when either {@code SimpleKeyGenerator}, {@code ComplexKeyGenerator}, diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java index 7aadac3d34079..db7901ee67c6d 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java @@ -50,6 +50,7 @@ import org.apache.hudi.common.model.HoodieRecordMerger; import org.apache.hudi.common.model.HoodieRecordPayload; import org.apache.hudi.common.model.HoodieTableType; +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; @@ -1772,8 +1773,39 @@ public int getSmallFileGroupCandidatesLimit() { return getInt(MERGE_SMALL_FILE_GROUP_CANDIDATES_LIMIT); } + /** + * @return true when every meta column is populated. + * + *

Derived from {@link #getMetaFieldsMode()} so that call sites still written against the + * deprecated {@code hoodie.populate.meta.fields} boolean observe the same answer as the enum: + * only {@link MetaFieldsMode#ALL} populates every meta column. + */ public boolean populateMetaFields() { - return getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS); + return getMetaFieldsMode().toLegacyPopulateMetaFields(); + } + + /** + * @return the {@link MetaFieldsMode} resolved from the write config. + * {@code hoodie.meta.fields.mode} is the source of truth; configs written before that property + * existed fall back to {@link MetaFieldsMode#ALL} or {@link MetaFieldsMode#NONE} based on the + * deprecated {@code hoodie.populate.meta.fields} boolean. + */ + public MetaFieldsMode getMetaFieldsMode() { + return MetaFieldsMode.resolve(this); + } + + /** + * @return true when {@code _hoodie_commit_time} is physically populated on every row. + */ + public boolean isCommitTimePopulated() { + return getMetaFieldsMode().isCommitTimePopulated(); + } + + /** + * @return true when {@code _hoodie_file_name} is physically populated on every row. + */ + public boolean isFileNamePopulated() { + return getMetaFieldsMode().isFileNamePopulated(); } /** @@ -3586,11 +3618,31 @@ public Builder withCanIgnorePostCommitFailures(boolean canIgnorePostCommitFailur return this; } + /** + * @deprecated since 1.3.0, use {@link #withMetaFieldsMode(MetaFieldsMode)} instead + * ({@code true} maps to {@link MetaFieldsMode#ALL}, {@code false} to {@link MetaFieldsMode#NONE}). + */ + @Deprecated public Builder withPopulateMetaFields(boolean populateMetaFields) { writeConfig.setValue(HoodieTableConfig.POPULATE_META_FIELDS, Boolean.toString(populateMetaFields)); return this; } + public Builder withMetaFieldsMode(MetaFieldsMode metaFieldsMode) { + // Leaving the mode unset defers to the deprecated populate.meta.fields boolean. Setting it + // also rewrites that boolean from the mode, so the two can never disagree — a config carrying + // a selective mode alongside populate.meta.fields=true would otherwise create a table whose + // hoodie.properties misleads pre-1.3.0 readers into treating it as ALL. + if (metaFieldsMode == null) { + writeConfig.setValue(HoodieTableConfig.META_FIELDS_MODE, ""); + } else { + writeConfig.setValue(HoodieTableConfig.META_FIELDS_MODE, metaFieldsMode.name()); + writeConfig.setValue(HoodieTableConfig.POPULATE_META_FIELDS, + Boolean.toString(metaFieldsMode.toLegacyPopulateMetaFields())); + } + return this; + } + public Builder withAllowOperationMetadataField(boolean allowOperationMetadataField) { writeConfig.setValue(ALLOW_OPERATION_METADATA_FIELD, Boolean.toString(allowOperationMetadataField)); return this; @@ -3883,6 +3935,27 @@ private void validate() { checkArgument(ttlStatsMaxParallelism > 0, String.format("%s must be positive, but was %d", HoodieTTLConfig.STATS_MAX_PARALLELISM.key(), ttlStatsMaxParallelism)); + + // hoodie.meta.fields.mode is the source of truth for meta-column population; the deprecated + // populate.meta.fields boolean is consulted only when the mode is absent. There is therefore + // no ambiguous combination to reject here — MetaFieldsMode.resolve throws on unrecognized + // values. + MetaFieldsMode metaFieldsMode = writeConfig.getMetaFieldsMode(); + // Selective meta-field modes are CoW-only in this release. MoR log-write path does not yet + // respect the mode, which would silently produce log records with null meta columns. + boolean isSelective = metaFieldsMode != MetaFieldsMode.ALL && metaFieldsMode != MetaFieldsMode.NONE; + checkArgument(!(writeConfig.getTableType() == HoodieTableType.MERGE_ON_READ && isSelective), + String.format("%s=%s is currently supported for COPY_ON_WRITE tables only. MoR support is a follow-up. " + + "For MoR use %s=ALL or %s=NONE.", + HoodieTableConfig.META_FIELDS_MODE.key(), metaFieldsMode, + HoodieTableConfig.META_FIELDS_MODE.key(), HoodieTableConfig.META_FIELDS_MODE.key())); + // Selective meta-field modes are wired only for the Spark writer path in this release. Flink + // RowData / Java-client writers ignore the mode and would silently produce NONE-mode output. + checkArgument(!(engineType != EngineType.SPARK && isSelective), + String.format("%s=%s is currently supported for the Spark writer only. Support for engine=%s is a follow-up. " + + "Use %s=ALL or %s=NONE.", + HoodieTableConfig.META_FIELDS_MODE.key(), metaFieldsMode, engineType, + HoodieTableConfig.META_FIELDS_MODE.key(), HoodieTableConfig.META_FIELDS_MODE.key())); } public HoodieWriteConfig build() { diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/index/partitionstats/PartitionStatsIndexer.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/index/partitionstats/PartitionStatsIndexer.java index fa66d0fa12e34..0dd0ae9b902c6 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/index/partitionstats/PartitionStatsIndexer.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/index/partitionstats/PartitionStatsIndexer.java @@ -28,6 +28,7 @@ import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieReplaceCommitMetadata; import org.apache.hudi.common.model.HoodieWriteStat; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.common.model.WriteOperationType; import org.apache.hudi.common.schema.HoodieSchema; import org.apache.hudi.common.schema.HoodieSchemaUtils; @@ -141,7 +142,10 @@ public static HoodieData convertMetadataToPartitionStatsRecords(Ho ? Option.empty() : Option.of(HoodieSchema.parse(writerSchemaStr))); HoodieTableConfig tableConfig = dataMetaClient.getTableConfig(); - Option tableSchema = writerSchema.map(schema -> tableConfig.populateMetaFields() ? HoodieSchemaUtils.addMetadataFields(schema) : schema); + // Selective meta-fields modes write the meta columns as physical nullable columns, so they + // belong in the table schema whenever the mode populates any of them. + Option tableSchema = writerSchema.map(schema -> + tableConfig.getMetaFieldsMode() != MetaFieldsMode.NONE ? HoodieSchemaUtils.addMetadataFields(schema) : schema); if (tableSchema.isEmpty()) { return engineContext.emptyHoodieData(); diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/NineToTenUpgradeHandler.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/NineToTenUpgradeHandler.java index 97fdea90036fa..ea7224991f512 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/NineToTenUpgradeHandler.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/NineToTenUpgradeHandler.java @@ -18,13 +18,25 @@ package org.apache.hudi.table.upgrade; +import org.apache.hudi.common.config.ConfigProperty; import org.apache.hudi.common.engine.HoodieEngineContext; +import org.apache.hudi.common.model.MetaFieldsMode; +import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.config.HoodieWriteConfig; +import java.util.Collections; +import java.util.Map; + /** * Version 10 enables native log format by default for new writes. Existing version 9 * inline log files remain readable by version 10 readers, so there is no table metadata * rewrite required for the upgrade. + * + *

Version 10 also introduced {@code hoodie.meta.fields.mode}. Version 9 tables predate it and + * resolve to {@code ALL} / {@code NONE} from the deprecated {@code hoodie.populate.meta.fields} + * boolean. The upgrade records that derived value explicitly so upgraded tables describe their + * meta-field layout the same way newly created version 10 tables do, rather than depending on the + * legacy fallback. Behavior is unchanged either way — this only makes the on-disk state explicit. */ public class NineToTenUpgradeHandler implements UpgradeHandler { @@ -34,6 +46,12 @@ public UpgradeDowngrade.TableConfigChangeSet upgrade( HoodieEngineContext context, String instantTime, SupportsUpgradeDowngrade upgradeDowngradeHelper) { - return new UpgradeDowngrade.TableConfigChangeSet(); + HoodieTableConfig tableConfig = + upgradeDowngradeHelper.getTable(config, context).getMetaClient().getTableConfig(); + // Resolves from the legacy boolean for a version 9 table, since the mode property is absent. + MetaFieldsMode metaFieldsMode = tableConfig.getMetaFieldsMode(); + Map propertiesToUpdate = Collections.singletonMap( + HoodieTableConfig.META_FIELDS_MODE, metaFieldsMode.name()); + return new UpgradeDowngrade.TableConfigChangeSet(propertiesToUpdate, Collections.emptySet()); } } diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/TenToNineDowngradeHandler.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/TenToNineDowngradeHandler.java index 393dd6986d822..dbeb53d9b07a9 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/TenToNineDowngradeHandler.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/TenToNineDowngradeHandler.java @@ -18,25 +18,61 @@ package org.apache.hudi.table.upgrade; +import org.apache.hudi.common.config.ConfigProperty; import org.apache.hudi.common.engine.HoodieEngineContext; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.config.HoodieWriteConfig; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import java.util.Collections; +import java.util.HashSet; +import java.util.Set; /** * Version 10 writes native log files by default. Downgrading to version 9 requires * full compaction of native data/delete logs before the downgrade completes. + * + *

Version 10 also introduced {@code hoodie.meta.fields.mode}. Version 9 does not understand it, + * so the property is dropped here while {@code hoodie.populate.meta.fields} is left exactly as it + * stands — {@code ALL} and {@code NONE} tables round-trip unchanged because those are precisely the + * two states the legacy boolean can express. Selective modes cannot be expressed in version 9, so + * the table degrades to what its legacy boolean says (which is {@code false}, i.e. NONE) and we warn. */ public class TenToNineDowngradeHandler implements DowngradeHandler { + + private static final Logger LOG = LoggerFactory.getLogger(TenToNineDowngradeHandler.class); + @Override public UpgradeDowngrade.TableConfigChangeSet downgrade( HoodieWriteConfig config, HoodieEngineContext context, String instantTime, SupportsUpgradeDowngrade upgradeDowngradeHelper) { + Set propertiesToDelete = new HashSet<>(); + propertiesToDelete.add(HoodieTableConfig.TABLE_STORAGE_LAYOUT); + + // The warning is best-effort: dropping the property is what matters, and the helper is not + // always available (some callers drive the change set directly). + MetaFieldsMode metaFieldsMode = upgradeDowngradeHelper == null + ? MetaFieldsMode.ALL + : upgradeDowngradeHelper.getTable(config, context).getMetaClient().getTableConfig().getMetaFieldsMode(); + if (metaFieldsMode != MetaFieldsMode.ALL && metaFieldsMode != MetaFieldsMode.NONE) { + LOG.warn("Table is using {}={}, which table version 9 cannot express. The property is being " + + "removed and the table will behave as {}=false (no meta columns) to version 9 readers. " + + "Already-written files keep their populated meta columns, but incremental queries that " + + "relied on {} will stop returning rows. Recreate the table if you need that behavior back.", + HoodieTableConfig.META_FIELDS_MODE.key(), metaFieldsMode, + HoodieTableConfig.POPULATE_META_FIELDS.key(), metaFieldsMode); + } + // hoodie.populate.meta.fields is deliberately left untouched: whatever the table recorded before + // the downgrade stays, so ALL and NONE tables are bit-identical afterwards. + propertiesToDelete.add(HoodieTableConfig.META_FIELDS_MODE); + return new UpgradeDowngrade.TableConfigChangeSet( Collections.emptyMap(), - Collections.singleton(HoodieTableConfig.TABLE_STORAGE_LAYOUT)); + propertiesToDelete); } } 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 4396e061a49ed..c318f5286a5e4 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 @@ -24,6 +24,7 @@ import org.apache.hudi.common.engine.HoodieLocalEngineContext; import org.apache.hudi.common.model.HoodieCommitMetadata; import org.apache.hudi.common.model.HoodieTimelineTimeZone; +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; @@ -43,6 +44,7 @@ import org.apache.hudi.config.HoodieLockConfig; import org.apache.hudi.config.HoodieWriteConfig; import org.apache.hudi.core.transaction.lock.InProcessLockProvider; +import org.apache.hudi.exception.HoodieException; import org.apache.hudi.index.HoodieIndex; import org.apache.hudi.index.HoodieSimpleIndex; import org.apache.hudi.keygen.ComplexAvroKeyGenerator; @@ -74,6 +76,8 @@ import static org.apache.hudi.common.testutils.HoodieTestUtils.getDefaultStorageConf; import static org.apache.hudi.testutils.Assertions.assertComplexKeyGeneratorValidationThrows; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.RETURNS_DEEP_STUBS; import static org.mockito.Mockito.mock; @@ -82,6 +86,92 @@ class TestBaseHoodieWriteClient extends HoodieCommonTestHarness { + private static HoodieTableConfig tableConfigWithMode(MetaFieldsMode mode) { + HoodieTableConfig tableConfig = new HoodieTableConfig(); + tableConfig.setValue(HoodieTableConfig.META_FIELDS_MODE, mode.name()); + tableConfig.setValue(HoodieTableConfig.POPULATE_META_FIELDS, + Boolean.toString(mode.toLegacyPopulateMetaFields())); + tableConfig.setValue(HoodieTableConfig.VERSION, + String.valueOf(HoodieTableVersion.current().versionCode())); + return tableConfig; + } + + private static BaseHoodieWriteClient validatorClient(HoodieWriteConfig writeConfig) { + return new TestWriteClient(writeConfig, mock(HoodieTable.class), Option.empty(), + mock(BaseHoodieTableServiceClient.class)); + } + + @Test + void validateAgainstTablePropertiesRejectsMetaFieldsModeMismatch() throws IOException { + initMetaClient(); + // A writer that explicitly asks for NONE against a COMMIT_TIME_ONLY table: both legacy booleans + // are false, so a boolean-only check passes and the writer goes on to produce null commit times + // while the table still advertises COMMIT_TIME_ONLY. + HoodieWriteConfig noneWriteConfig = HoodieWriteConfig.newBuilder() + .withPath(basePath) + .withMetaFieldsMode(MetaFieldsMode.NONE) + .build(); + assertEquals(MetaFieldsMode.NONE, noneWriteConfig.getMetaFieldsMode()); + + HoodieTableConfig commitTimeOnlyTable = tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY); + assertFalse(commitTimeOnlyTable.populateMetaFields(), + "precondition: both legacy booleans are false, so only the enum comparison can catch this"); + + HoodieException ex = assertThrows(HoodieException.class, () -> + validatorClient(noneWriteConfig).validateAgainstTableProperties(commitTimeOnlyTable, noneWriteConfig)); + 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"), + "error must name both modes: " + ex.getMessage()); + } + + @Test + void validateAgainstTablePropertiesAllowsUnstatedWriterToNarrow() throws IOException { + initMetaClient(); + // Long-standing behavior: a writer that sets only populate.meta.fields=false (never naming a + // mode) resolves to NONE, and must still be able to write to an ALL table. Many callers build + // a write config without restating the table's meta-field settings. + HoodieWriteConfig unstated = HoodieWriteConfig.newBuilder() + .withPath(basePath) + .withPopulateMetaFields(false) + .build(); + assertEquals(MetaFieldsMode.NONE, unstated.getMetaFieldsMode()); + validatorClient(unstated) + .validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.ALL), unstated); + } + + @Test + void validateAgainstTablePropertiesRejectsWideningEvenWhenUnstated() throws IOException { + initMetaClient(); + // Widening is rejected regardless of whether the writer named a mode: a default writer resolves + // to ALL, which would claim meta columns a NONE table never wrote. + HoodieWriteConfig defaultWriteConfig = HoodieWriteConfig.newBuilder().withPath(basePath).build(); + assertEquals(MetaFieldsMode.ALL, defaultWriteConfig.getMetaFieldsMode()); + + HoodieException ex = assertThrows(HoodieException.class, () -> + validatorClient(defaultWriteConfig) + .validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.NONE), defaultWriteConfig)); + assertTrue(ex.getMessage().contains("cannot be widened"), ex.getMessage()); + } + + @Test + void validateAgainstTablePropertiesAcceptsMatchingMetaFieldsMode() throws IOException { + initMetaClient(); + // Default writer and default table both resolve to ALL — the overwhelmingly common case. + HoodieWriteConfig defaultWriteConfig = HoodieWriteConfig.newBuilder().withPath(basePath).build(); + assertEquals(MetaFieldsMode.ALL, defaultWriteConfig.getMetaFieldsMode()); + validatorClient(defaultWriteConfig) + .validateAgainstTableProperties(tableConfigWithMode(MetaFieldsMode.ALL), defaultWriteConfig); + + // And a selective writer against a table recorded with the same mode. + HoodieWriteConfig selectiveWriteConfig = HoodieWriteConfig.newBuilder() + .withPath(basePath) + .withMetaFieldsMode(MetaFieldsMode.COMMIT_TIME_ONLY) + .build(); + validatorClient(selectiveWriteConfig).validateAgainstTableProperties( + tableConfigWithMode(MetaFieldsMode.COMMIT_TIME_ONLY), selectiveWriteConfig); + } + @Test void startCommitWillRollbackFailedWritesInEagerMode() throws IOException { initMetaClient(); diff --git a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfigMetaFieldsMode.java b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfigMetaFieldsMode.java new file mode 100644 index 0000000000000..aff79baa3de3e --- /dev/null +++ b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfigMetaFieldsMode.java @@ -0,0 +1,189 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hudi.config; + +import org.apache.hudi.common.model.HoodieTableType; +import org.apache.hudi.common.model.MetaFieldsMode; +import org.apache.hudi.common.table.HoodieTableConfig; + +import org.junit.jupiter.api.Test; + +import java.util.Properties; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Validates the writer-side accessors and validation guards for the meta-field-population modes + * on {@link HoodieWriteConfig}. Companion test for the {@link HoodieTableConfig} accessors lives + * in {@code TestHoodieMetaFieldsMode}; this test covers the writer-builder surface and the + * cross-flag validation that runs at {@code build()} time. + */ +class TestHoodieWriteConfigMetaFieldsMode { + + private static HoodieWriteConfig.Builder baseBuilder() { + return HoodieWriteConfig.newBuilder().withPath("file:///tmp/test_hudi_meta_fields_mode"); + } + + private static Properties mergeOnReadProps() { + Properties props = new Properties(); + props.setProperty(HoodieTableConfig.TYPE.key(), HoodieTableType.MERGE_ON_READ.name()); + return props; + } + + @Test + void defaultsToAllMode() { + HoodieWriteConfig cfg = baseBuilder().build(); + assertTrue(cfg.populateMetaFields()); + assertEquals(MetaFieldsMode.ALL, cfg.getMetaFieldsMode()); + assertTrue(cfg.isCommitTimePopulated()); + assertTrue(cfg.isFileNamePopulated()); + } + + @Test + void explicitNoneModeBuilds() { + HoodieWriteConfig cfg = baseBuilder().withPopulateMetaFields(false).build(); + assertFalse(cfg.populateMetaFields()); + assertEquals(MetaFieldsMode.NONE, cfg.getMetaFieldsMode()); + assertFalse(cfg.isCommitTimePopulated()); + assertFalse(cfg.isFileNamePopulated()); + } + + @Test + void commitTimeOnlyModeBuilds() { + HoodieWriteConfig cfg = baseBuilder() + .withPopulateMetaFields(false) + .withMetaFieldsMode(MetaFieldsMode.COMMIT_TIME_ONLY) + .build(); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, cfg.getMetaFieldsMode()); + assertTrue(cfg.isCommitTimePopulated()); + assertFalse(cfg.isFileNamePopulated()); + } + + @Test + void fileNameOnlyModeBuilds() { + HoodieWriteConfig cfg = baseBuilder() + .withPopulateMetaFields(false) + .withMetaFieldsMode(MetaFieldsMode.FILE_NAME_ONLY) + .build(); + assertEquals(MetaFieldsMode.FILE_NAME_ONLY, cfg.getMetaFieldsMode()); + assertFalse(cfg.isCommitTimePopulated()); + assertTrue(cfg.isFileNamePopulated()); + } + + @Test + void commitTimeAndFileNameCombinationBuilds() { + HoodieWriteConfig cfg = baseBuilder() + .withPopulateMetaFields(false) + .withMetaFieldsMode(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME) + .build(); + assertEquals(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME, cfg.getMetaFieldsMode()); + assertTrue(cfg.isCommitTimePopulated()); + assertTrue(cfg.isFileNamePopulated()); + } + + @Test + void modeWinsOverLegacyPopulateMetaFields() { + // The two properties no longer compose: an explicit mode is the answer regardless of the + // deprecated boolean, so this combination is accepted rather than rejected. + HoodieWriteConfig cfg = baseBuilder() + .withPopulateMetaFields(true) + .withMetaFieldsMode(MetaFieldsMode.COMMIT_TIME_ONLY) + .build(); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, cfg.getMetaFieldsMode()); + assertTrue(cfg.isCommitTimePopulated()); + assertFalse(cfg.isFileNamePopulated()); + // populateMetaFields() is derived from the mode — only ALL reports true. + assertFalse(cfg.populateMetaFields()); + // The raw legacy property is rewritten too, so a config handed to table creation cannot carry + // a boolean that contradicts the mode. + assertEquals("false", cfg.getStringOrDefault(HoodieTableConfig.POPULATE_META_FIELDS)); + } + + @Test + void legacyBooleanSetAfterModeStillResolvesFromMode() { + // Builder call order must not change the outcome: getMetaFieldsMode() reads the mode property + // first, so a later withPopulateMetaFields(...) cannot silently widen a selective table. + HoodieWriteConfig cfg = baseBuilder() + .withMetaFieldsMode(MetaFieldsMode.COMMIT_TIME_ONLY) + .withPopulateMetaFields(true) + .build(); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, cfg.getMetaFieldsMode()); + assertFalse(cfg.populateMetaFields()); + } + + @Test + void explicitAllModeOverridesLegacyFalse() { + HoodieWriteConfig cfg = baseBuilder() + .withPopulateMetaFields(false) + .withMetaFieldsMode(MetaFieldsMode.ALL) + .build(); + assertEquals(MetaFieldsMode.ALL, cfg.getMetaFieldsMode()); + assertTrue(cfg.populateMetaFields()); + } + + @Test + void noneModeWithExplicitBuildIsStillNone() { + HoodieWriteConfig cfg = baseBuilder() + .withPopulateMetaFields(false) + .withMetaFieldsMode(MetaFieldsMode.NONE) + .build(); + assertEquals(MetaFieldsMode.NONE, cfg.getMetaFieldsMode()); + assertFalse(cfg.isCommitTimePopulated()); + assertFalse(cfg.isFileNamePopulated()); + } + + @Test + void legacyBooleanIsUsedWhenModeIsAbsent() { + // Backward compat: tables written before hoodie.meta.fields.mode existed keep their behavior. + HoodieWriteConfig allCfg = baseBuilder().withPopulateMetaFields(true).build(); + assertEquals(MetaFieldsMode.ALL, allCfg.getMetaFieldsMode()); + + HoodieWriteConfig noneCfg = baseBuilder().withPopulateMetaFields(false).build(); + assertEquals(MetaFieldsMode.NONE, noneCfg.getMetaFieldsMode()); + } + + @Test + void rejectsSelectiveModeOnMergeOnRead() { + // Selective modes are CoW-only until the MoR log-write path honors them. + HoodieWriteConfig.Builder builder = baseBuilder() + .withMetaFieldsMode(MetaFieldsMode.COMMIT_TIME_ONLY) + .withProps(mergeOnReadProps()); + IllegalArgumentException ex = assertThrows(IllegalArgumentException.class, builder::build); + assertTrue(ex.getMessage().contains("hoodie.meta.fields.mode"), + "exception must name the mode property: " + ex.getMessage()); + assertTrue(ex.getMessage().contains("COPY_ON_WRITE"), + "exception must explain the CoW-only restriction: " + ex.getMessage()); + } + + @Test + void allowsAllAndNoneOnMergeOnRead() { + // Only the selective modes are restricted — the two legacy-equivalent modes stay available. + assertEquals(MetaFieldsMode.ALL, baseBuilder() + .withMetaFieldsMode(MetaFieldsMode.ALL) + .withProps(mergeOnReadProps()) + .build().getMetaFieldsMode()); + assertEquals(MetaFieldsMode.NONE, baseBuilder() + .withMetaFieldsMode(MetaFieldsMode.NONE) + .withProps(mergeOnReadProps()) + .build().getMetaFieldsMode()); + } +} diff --git a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestNineToTenUpgradeHandler.java b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestNineToTenUpgradeHandler.java new file mode 100644 index 0000000000000..b4774cad863fc --- /dev/null +++ b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestNineToTenUpgradeHandler.java @@ -0,0 +1,85 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hudi.table.upgrade; + +import org.apache.hudi.common.engine.HoodieEngineContext; +import org.apache.hudi.common.model.MetaFieldsMode; +import org.apache.hudi.common.table.HoodieTableConfig; +import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.config.HoodieWriteConfig; +import org.apache.hudi.table.HoodieTable; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.RETURNS_DEEP_STUBS; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * Version 9 tables predate {@code hoodie.meta.fields.mode}, so the upgrade records the value + * derived from the deprecated {@code hoodie.populate.meta.fields} boolean. This makes an upgraded + * table describe its meta-field layout the same way a freshly created version 10 table does, + * instead of relying on the legacy fallback at every read. + */ +class TestNineToTenUpgradeHandler { + + private static SupportsUpgradeDowngrade helperFor(MetaFieldsMode resolvedMode) { + HoodieTable table = mock(HoodieTable.class, RETURNS_DEEP_STUBS); + HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class, RETURNS_DEEP_STUBS); + HoodieTableConfig tableConfig = mock(HoodieTableConfig.class); + when(tableConfig.getMetaFieldsMode()).thenReturn(resolvedMode); + when(metaClient.getTableConfig()).thenReturn(tableConfig); + when(table.getMetaClient()).thenReturn(metaClient); + + SupportsUpgradeDowngrade helper = mock(SupportsUpgradeDowngrade.class); + when(helper.getTable(org.mockito.ArgumentMatchers.any(HoodieWriteConfig.class), + org.mockito.ArgumentMatchers.any(HoodieEngineContext.class))).thenReturn(table); + return helper; + } + + @ParameterizedTest + @CsvSource({"ALL", "NONE"}) + void upgradeRecordsTheModeDerivedFromTheLegacyBoolean(String modeName) { + MetaFieldsMode expected = MetaFieldsMode.valueOf(modeName); + UpgradeDowngrade.TableConfigChangeSet changeSet = new NineToTenUpgradeHandler().upgrade( + mock(HoodieWriteConfig.class), mock(HoodieEngineContext.class), "001", helperFor(expected)); + + assertTrue(changeSet.propertiesToDelete().isEmpty()); + assertEquals(1, changeSet.propertiesToUpdate().size()); + assertEquals(expected.name(), + changeSet.propertiesToUpdate().get(HoodieTableConfig.META_FIELDS_MODE)); + } + + @Test + void upgradeLeavesTheLegacyBooleanAlone() { + // The boolean stays authoritative for any reader that has not learned about the mode yet, and + // the two must agree — so the upgrade only adds the mode, never rewrites populate.meta.fields. + UpgradeDowngrade.TableConfigChangeSet changeSet = new NineToTenUpgradeHandler().upgrade( + mock(HoodieWriteConfig.class), mock(HoodieEngineContext.class), "001", + helperFor(MetaFieldsMode.NONE)); + + assertTrue(changeSet.propertiesToDelete().isEmpty()); + assertEquals(1, changeSet.propertiesToUpdate().size()); + assertTrue(changeSet.propertiesToUpdate().containsKey(HoodieTableConfig.META_FIELDS_MODE)); + } +} diff --git a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestTenToNineDowngradeHandler.java b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestTenToNineDowngradeHandler.java index 12faa97369303..b9383f6d2be86 100644 --- a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestTenToNineDowngradeHandler.java +++ b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/upgrade/TestTenToNineDowngradeHandler.java @@ -24,18 +24,25 @@ import org.junit.jupiter.api.Test; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; class TestTenToNineDowngradeHandler { @Test - void testDowngradeRemovesStorageLayoutOnly() { + void testDowngradeRemovesStorageLayoutAndMetaFieldsMode() { UpgradeDowngrade.TableConfigChangeSet changeSet = new TenToNineDowngradeHandler().downgrade(null, null, null, null); assertTrue(changeSet.propertiesToUpdate().isEmpty()); - assertEquals(1, changeSet.propertiesToDelete().size()); + assertEquals(2, changeSet.propertiesToDelete().size()); assertTrue(changeSet.propertiesToDelete().contains(HoodieTableConfig.TABLE_STORAGE_LAYOUT)); + // Version 9 does not understand hoodie.meta.fields.mode, so it is dropped... + assertTrue(changeSet.propertiesToDelete().contains(HoodieTableConfig.META_FIELDS_MODE)); + // ...while hoodie.populate.meta.fields is deliberately left in place, so ALL and NONE tables + // round-trip unchanged — those are exactly the two states the legacy boolean can express. + assertFalse(changeSet.propertiesToDelete().contains(HoodieTableConfig.POPULATE_META_FIELDS)); + assertFalse(changeSet.propertiesToUpdate().containsKey(HoodieTableConfig.POPULATE_META_FIELDS)); } @Test @@ -44,7 +51,8 @@ void testTenToNineDowngradeRouteIsSupported() { new UpgradeDowngrade(null, null, null, null) .downgrade(HoodieTableVersion.TEN, HoodieTableVersion.NINE, "001"); - assertEquals(1, changeSet.propertiesToDelete().size()); + assertEquals(2, changeSet.propertiesToDelete().size()); assertTrue(changeSet.propertiesToDelete().contains(HoodieTableConfig.TABLE_STORAGE_LAYOUT)); + assertTrue(changeSet.propertiesToDelete().contains(HoodieTableConfig.META_FIELDS_MODE)); } } diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkMetadataWriterUtils.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkMetadataWriterUtils.java index 5a4b354e868d6..c3d292fdb6423 100644 --- a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkMetadataWriterUtils.java +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkMetadataWriterUtils.java @@ -42,6 +42,7 @@ import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieReplaceCommitMetadata; import org.apache.hudi.common.model.HoodieWriteStat; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.common.schema.HoodieSchema; import org.apache.hudi.common.schema.HoodieSchemaUtils; import org.apache.hudi.common.table.HoodieTableConfig; @@ -391,7 +392,10 @@ public static HoodiePairData> ? Option.empty() : Option.of(HoodieSchema.parse(writerSchemaStr))); HoodieTableConfig tableConfig = dataMetaClient.getTableConfig(); - HoodieSchema tableSchema = writerSchema.map(schema -> tableConfig.populateMetaFields() ? HoodieSchemaUtils.addMetadataFields(schema) : schema) + // Selective meta-fields modes write the meta columns as physical nullable columns, so they + // belong in the table schema whenever the mode populates any of them. + HoodieSchema tableSchema = writerSchema.map(schema -> + tableConfig.getMetaFieldsMode() != MetaFieldsMode.NONE ? HoodieSchemaUtils.addMetadataFields(schema) : schema) .orElseThrow(() -> new IllegalStateException(String.format("Expected writer schema in commit metadata %s", commitMetadata))); List> columnsToIndexSchemaMap = columnsToIndex.stream() .map(columnToIndex -> HoodieSchemaUtils.getNestedField(tableSchema, columnToIndex)) diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkValidatorUtils.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkValidatorUtils.java index 233f622beff68..b57dc68ef46a5 100644 --- a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkValidatorUtils.java +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkValidatorUtils.java @@ -27,6 +27,7 @@ import org.apache.hudi.common.engine.HoodieEngineContext; import org.apache.hudi.common.model.BaseFile; import org.apache.hudi.common.model.HoodieWriteStat; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.common.schema.HoodieSchema; import org.apache.hudi.common.schema.HoodieSchemaUtils; import org.apache.hudi.common.table.TableSchemaResolver; @@ -197,10 +198,13 @@ public static Dataset readRecordsForBaseFiles(SQLContext sqlContext, List injectedConfigs = HoodieParquetConfigInjector.applyConfigInjector(path, storage.getConf(), config); StorageConfiguration storageConfiguration = injectedConfigs.getLeft(); @@ -80,7 +82,7 @@ protected HoodieFileWriter newParquetFileWriter( hoodieConfig.getBooleanOrDefault(HoodieStorageConfig.PARQUET_DICTIONARY_ENABLED)); parquetConfig.getHadoopConf().addResource(writeSupport.getHadoopConf()); - return new HoodieSparkParquetWriter(path, parquetConfig, instantTime, taskContextSupplier, populateMetaFields); + return new HoodieSparkParquetWriter(path, parquetConfig, instantTime, taskContextSupplier, metaFieldsMode); } protected HoodieFileWriter newParquetFileWriter(OutputStream outputStream, HoodieConfig config, @@ -119,7 +121,12 @@ protected HoodieFileWriter newOrcFileWriter(String instantTime, StoragePath path protected HoodieFileWriter newLanceFileWriter(String instantTime, StoragePath path, HoodieConfig config, HoodieSchema schema, TaskContextSupplier taskContextSupplier) throws IOException { HoodieSparkLanceWriter.validateNoVariantColumns(schema); - boolean populateMetaFields = config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS); + // Resolve through hoodie.meta.fields.mode rather than the deprecated boolean — see the parquet + // path above. Lance does not yet populate meta columns selectively, so a selective mode is + // treated as "record key not populated" (no bloom filter, no meta stamping). + boolean populateMetaFields = org.apache.hudi.common.model.MetaFieldsMode.resolve( + config.getStringOrDefault(HoodieTableConfig.META_FIELDS_MODE), + config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS)).isRecordKeyPopulated(); StructType structType = HoodieInternalRowUtils.getCachedSchema(schema); boolean enableBloomFilter = enableBloomFilter(populateMetaFields, config); Option bloomFilter = enableBloomFilter ? Option.of(createBloomFilter(config)) : Option.empty(); @@ -144,7 +151,12 @@ protected HoodieFileWriter newLanceFileWriter(String instantTime, StoragePath pa @Override protected HoodieFileWriter newVortexFileWriter(String instantTime, StoragePath path, HoodieConfig config, HoodieSchema schema, TaskContextSupplier taskContextSupplier) throws IOException { - boolean populateMetaFields = config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS); + // Resolve through hoodie.meta.fields.mode rather than the deprecated boolean — see the parquet + // path above. Vortex does not yet populate meta columns selectively, so a selective mode is + // treated as "record key not populated" (no bloom filter, no meta stamping). + boolean populateMetaFields = org.apache.hudi.common.model.MetaFieldsMode.resolve( + config.getStringOrDefault(HoodieTableConfig.META_FIELDS_MODE), + config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS)).isRecordKeyPopulated(); StructType structType = HoodieInternalRowUtils.getCachedSchema(schema); boolean enableBloomFilter = enableBloomFilter(populateMetaFields, config); Option bloomFilter = enableBloomFilter ? Option.of(createBloomFilter(config)) : Option.empty(); diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkParquetWriter.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkParquetWriter.java index f2ed3831d492a..9480d85f9ea34 100644 --- a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkParquetWriter.java +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkParquetWriter.java @@ -21,6 +21,7 @@ import org.apache.hudi.common.engine.TaskContextSupplier; import org.apache.hudi.common.model.HoodieKey; import org.apache.hudi.common.model.HoodieRecord; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.io.hadoop.HoodieBaseParquetWriter; import org.apache.hudi.io.storage.row.HoodieRowParquetConfig; import org.apache.hudi.io.storage.row.HoodieRowParquetWriteSupport; @@ -44,7 +45,7 @@ public class HoodieSparkParquetWriter extends HoodieBaseParquetWriter { Integer partitionId = taskContextSupplier.getPartitionIdSupplier().get(); return HoodieRecord.generateSequenceId(instantTime, partitionId, recordIndex); @@ -68,21 +69,35 @@ public HoodieSparkParquetWriter(StoragePath file, @Override public void writeRowWithMetadata(HoodieKey key, InternalRow row) throws IOException { - if (populateMetaFields) { - UTF8String recordKey = UTF8String.fromString(key.getRecordKey()); - updateRecordMetadata(row, recordKey, key.getPartitionPath(), getWrittenRecordCount()); - - super.write(row); - writeSupport.add(recordKey); - } else { - super.write(row); + switch (metaFieldsMode) { + case ALL: + UTF8String recordKey = UTF8String.fromString(key.getRecordKey()); + updateRecordMetadata(row, recordKey, key.getPartitionPath(), getWrittenRecordCount()); + super.write(row); + writeSupport.add(recordKey); + break; + case NONE: + super.write(row); + break; + default: + // Selective mode — populate only the opted-in columns. Record-key column stays null, so + // we do NOT register the record key with the write support (bloom filter / RLI hooks are + // meaningless without the record-key column). + if (metaFieldsMode.isCommitTimePopulated()) { + row.update(COMMIT_TIME_METADATA_FIELD.ordinal(), instantTime); + } + if (metaFieldsMode.isFileNamePopulated()) { + row.update(FILENAME_METADATA_FIELD.ordinal(), fileName); + } + super.write(row); + break; } } @Override public void writeRow(String recordKey, InternalRow row) throws IOException { super.write(row); - if (populateMetaFields) { + if (metaFieldsMode == MetaFieldsMode.ALL) { writeSupport.add(UTF8String.fromString(recordKey)); } } diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowCreateHandle.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowCreateHandle.java index c6d20eefb5823..aa6bdbc7ba3ae 100644 --- a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowCreateHandle.java +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowCreateHandle.java @@ -28,6 +28,7 @@ import org.apache.hudi.common.model.HoodieRecordLocation; import org.apache.hudi.common.model.HoodieWriteStat; import org.apache.hudi.common.model.IOType; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.common.util.HoodieTimer; import org.apache.hudi.common.util.Option; import org.apache.hudi.common.util.collection.Pair; @@ -66,7 +67,7 @@ public class HoodieRowCreateHandle implements Serializable { private final StoragePath path; private final String fileId; - private final boolean populateMetaFields; + private final MetaFieldsMode metaFieldsMode; private final UTF8String fileName; private final UTF8String commitTime; @@ -118,7 +119,7 @@ public HoodieRowCreateHandle(HoodieTable table, table.getBaseFileExtension()); this.path = makeNewPath(storage, partitionPath, fileName, writeConfig); - this.populateMetaFields = writeConfig.populateMetaFields(); + this.metaFieldsMode = writeConfig.getMetaFieldsMode(); this.fileName = UTF8String.fromString(path.getName()); this.commitTime = UTF8String.fromString(instantTime); this.seqIdGenerator = (id) -> HoodieRecord.generateSequenceId(instantTime, taskPartitionId, id); @@ -158,10 +159,48 @@ public HoodieRowCreateHandle(HoodieTable table, * @throws IOException */ public void write(InternalRow row) throws IOException { - if (populateMetaFields) { - writeRow(row); - } else { - writeRowNoMetaFields(row); + switch (metaFieldsMode) { + case ALL: + writeRow(row); + break; + case NONE: + writeRowNoMetaFields(row); + break; + default: + writeRowSelectiveMetaFields(row); + break; + } + } + + /** + * Selective meta-field write path: populate only the meta columns opted in via + * {@code hoodie.meta.fields.mode} — {@code _hoodie_commit_time} and/or {@code _hoodie_file_name}. + * The other meta columns stay null on disk. Record key is never populated in this path, so the + * record key is not registered with the write support (bloom filter / RLI hooks are meaningless + * without the record-key column). + */ + private void writeRowSelectiveMetaFields(InternalRow row) { + try { + UTF8String[] metaFields = new UTF8String[5]; + if (metaFieldsMode.isCommitTimePopulated()) { + metaFields[HoodieRecord.COMMIT_TIME_METADATA_FIELD_ORD] = shouldPreserveHoodieMetadata + ? row.getUTF8String(HoodieRecord.COMMIT_TIME_METADATA_FIELD_ORD) : commitTime; + } + if (metaFieldsMode.isFileNamePopulated()) { + // Always the file being written, never the source row's value — even when + // shouldPreserveHoodieMetadata is set. Preserving it during clustering would leave + // _hoodie_file_name pointing at a file that clustering just replaced. This matches the + // ALL path in writeRow below. + metaFields[HoodieRecord.FILENAME_META_FIELD_ORD] = fileName; + } + // The remaining meta columns stay null — Parquet stores nulls as definition-level flags + // (zero data bytes). + InternalRow updatedRow = SparkAdapterSupport$.MODULE$.sparkAdapter().createInternalRow(metaFields, row, true); + fileWriter.writeRow(updatedRow); + writeStatus.markSuccess((HoodieRecordDelegate) null, Option.empty()); + } catch (Exception e) { + writeStatus.setGlobalError(e); + throw new HoodieException("Exception thrown while writing spark InternalRows to file ", e); } } diff --git a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/HoodieDatasetBulkInsertHelper.scala b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/HoodieDatasetBulkInsertHelper.scala index 0ee2abd847a51..85c33e3f42741 100644 --- a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/HoodieDatasetBulkInsertHelper.scala +++ b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/HoodieDatasetBulkInsertHelper.scala @@ -137,7 +137,12 @@ object HoodieDatasetBulkInsertHelper // need access to the [[InternalRow]] and therefore can avoid the need // to dereference [[DataFrame]] into [[RDD]] val query = df.queryExecution.logical - val metaFieldsStubs = metaFields.map(f => Alias(Literal(UTF8String.EMPTY_UTF8, dataType = StringType), f.name)()) + // Nullable null stubs — the actual meta-column values are set downstream by + // HoodieRowCreateHandle.write based on hoodie.meta.fields.mode. Using a null literal (rather + // than an empty-string literal) guarantees the resulting StructField's nullable=true so the + // physical Parquet column is written as OPTIONAL and can hold nulls under selective / NONE + // modes. + val metaFieldsStubs = metaFields.map(f => Alias(Literal.create(null, StringType), f.name)()) val prependedQuery = Project(metaFieldsStubs ++ query.output, query) sparkAdapter.getUnsafeUtils.createDataFrameFrom(df.sparkSession, prependedQuery) diff --git a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroIndexedRecord.java b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroIndexedRecord.java index b2b3e99d345e8..ea08ac9e41d67 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroIndexedRecord.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroIndexedRecord.java @@ -22,6 +22,7 @@ import org.apache.hudi.common.avro.HoodieAvroUtils; import org.apache.hudi.common.avro.JoinedGenericRecord; import org.apache.hudi.common.schema.HoodieSchema; +import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.common.table.read.DeleteContext; import org.apache.hudi.common.util.ConfigUtils; import org.apache.hudi.common.util.Option; @@ -281,7 +282,14 @@ public HoodieRecord wrapIntoHoodieRecordPayloadWithKeyGen(HoodieSchema recordSch GenericRecord record = (GenericRecord) data; String key; String partition; - if (keyGen.isPresent() && !Boolean.parseBoolean(props.getOrDefault(POPULATE_META_FIELDS.key(), POPULATE_META_FIELDS.defaultValue().toString()).toString())) { + // Resolve via hoodie.meta.fields.mode — reading the deprecated boolean alone would report + // "populated" for a selective-mode table (whose _hoodie_record_key column is null), sending us + // down the meta-column branch below and NPE-ing on the null field. + boolean recordKeyPopulated = MetaFieldsMode.resolve( + props.getProperty(HoodieTableConfig.META_FIELDS_MODE.key()), + Boolean.parseBoolean(props.getOrDefault(POPULATE_META_FIELDS.key(), + POPULATE_META_FIELDS.defaultValue().toString()).toString())).isRecordKeyPopulated(); + if (keyGen.isPresent() && !recordKeyPopulated) { BaseKeyGenerator keyGeneratorOpt = keyGen.get(); key = keyGeneratorOpt.getRecordKey(record); partition = keyGeneratorOpt.getPartitionPath(record); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/model/MetaFieldsMode.java b/hudi-common/src/main/java/org/apache/hudi/common/model/MetaFieldsMode.java new file mode 100644 index 0000000000000..6467337f571e6 --- /dev/null +++ b/hudi-common/src/main/java/org/apache/hudi/common/model/MetaFieldsMode.java @@ -0,0 +1,191 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hudi.common.model; + +import org.apache.hudi.common.config.HoodieConfig; +import org.apache.hudi.common.table.HoodieTableConfig; +import org.apache.hudi.common.util.StringUtils; + +import java.util.Locale; + +/** + * Which of Hudi's meta columns are physically populated on disk. + * + *

Selective modes exist so that tables that opt out of the default {@code populate.meta.fields=true} + * can still keep the two columns that matter for downstream operations without paying for the other + * three: + * + *

+ * + *

The remaining three meta columns ({@code _hoodie_commit_seqno}, {@code _hoodie_record_key}, + * {@code _hoodie_partition_path}) are all-or-nothing — either populate every meta column ({@link #ALL}) + * or none of them beyond the two selectable ones. If you need any of the remaining columns, set + * {@code hoodie.populate.meta.fields=true}. + * + *

This enum is the single source of truth for meta-column population. The legacy boolean + * {@code hoodie.populate.meta.fields} is deprecated and consulted only when + * {@code hoodie.meta.fields.mode} is absent, so that tables written before the mode property + * existed keep their behavior: + * + *

+ * + *

On-disk representation: the enum {@link #name()} is persisted in {@code hoodie.properties} + * under the property {@code hoodie.meta.fields.mode}. + */ +public enum MetaFieldsMode { + /** + * All five Hudi meta columns are populated — today's default. + */ + ALL(true, true), + + /** + * No Hudi meta columns are populated. Incremental queries are unsupported. File-level pruning + * that depends on {@code _hoodie_file_name} is unsupported. + */ + NONE(false, false), + + /** + * Only {@code _hoodie_commit_time} is populated. Incremental queries remain functional; other + * meta columns stay null on disk. + */ + COMMIT_TIME_ONLY(true, false), + + /** + * Only {@code _hoodie_file_name} is populated. Useful for file-level lookups and debugging; + * incremental queries are unsupported. + */ + FILE_NAME_ONLY(false, true), + + /** + * Both {@code _hoodie_commit_time} and {@code _hoodie_file_name} are populated. + */ + COMMIT_TIME_AND_FILE_NAME(true, true); + + private final boolean commitTimePopulated; + private final boolean fileNamePopulated; + + MetaFieldsMode(boolean commitTimePopulated, boolean fileNamePopulated) { + this.commitTimePopulated = commitTimePopulated; + this.fileNamePopulated = fileNamePopulated; + } + + public boolean isCommitTimePopulated() { + return commitTimePopulated; + } + + public boolean isFileNamePopulated() { + return fileNamePopulated; + } + + /** + * @return true when all five meta columns are populated (i.e. this is {@link #ALL}). Selective + * modes never populate {@code _hoodie_record_key}, {@code _hoodie_partition_path}, or + * {@code _hoodie_commit_seqno}. + */ + public boolean isRecordKeyPopulated() { + return this == ALL; + } + + /** + * Resolve the effective mode. {@code hoodie.meta.fields.mode} is the source of truth; the + * deprecated {@code hoodie.populate.meta.fields} boolean is a fallback for tables written before + * the mode property existed. Precedence: + * + *

+ * + * @param rawMode raw {@code hoodie.meta.fields.mode} value; may be null or empty. + * @param legacyPopulateMetaFields value of the deprecated {@code hoodie.populate.meta.fields}. + * @throws IllegalArgumentException when the raw mode value does not match any enum value. This + * includes the pre-enum comma-separated format — callers that upgrade an old table must + * migrate the value through the hudi-cli. + */ + /** + * Resolve the effective mode from any {@link HoodieConfig} that may carry the two properties — + * a table config, a write config, or a bare config built from write options. Preferred over the + * two-argument overload: it keeps the property keys and the precedence rule in one place instead + * of repeating them at every call site. + */ + public static MetaFieldsMode resolve(HoodieConfig config) { + return resolve(config.getStringOrDefault(HoodieTableConfig.META_FIELDS_MODE), + config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS)); + } + + public static MetaFieldsMode resolve(String rawMode, boolean legacyPopulateMetaFields) { + if (StringUtils.isNullOrEmpty(rawMode)) { + return legacyPopulateMetaFields ? ALL : NONE; + } + return parse(rawMode); + } + + /** + * Parse a raw {@code hoodie.meta.fields.mode} value into an enum constant, with a message that + * lists the allowed values. Prefer this over {@link #valueOf(String)} for user-supplied input. + */ + public static MetaFieldsMode parse(String rawMode) { + try { + // Case-insensitive: users hand-editing hoodie.properties or passing write options should not + // have to match the enum's casing exactly. + return MetaFieldsMode.valueOf(rawMode.trim().toUpperCase(Locale.ROOT)); + } catch (IllegalArgumentException e) { + throw new IllegalArgumentException(String.format( + "Unsupported value '%s' for hoodie.meta.fields.mode. Allowed values: %s, %s, %s, %s, %s.", + rawMode, ALL, NONE, COMMIT_TIME_ONLY, FILE_NAME_ONLY, COMMIT_TIME_AND_FILE_NAME), e); + } + } + + /** + * @return the equivalent value of the deprecated {@code hoodie.populate.meta.fields} boolean, so + * that call sites not yet migrated to this enum keep observing consistent behavior. + */ + public boolean toLegacyPopulateMetaFields() { + return this == ALL; + } + + /** + * @return true when this mode populates at least one meta column that {@code other} does not. + * + *

Meta-field population is a physical-storage decision baked into files at write time, so it + * can never be widened for an existing table: earlier commits would be missing columns that later + * commits have, and readers cannot tell the two apart. Every transition that adds a column is + * therefore rejected — {@code NONE -> COMMIT_TIME_ONLY} and + * {@code FILE_NAME_ONLY -> COMMIT_TIME_AND_FILE_NAME} just as much as {@code NONE -> ALL}. + * + *

Narrowing is not flagged here: writing fewer meta columns than the table advertises cannot + * make a reader believe in data that is absent, and it is long-standing behavior for a writer to + * resolve to {@link #NONE} against an {@link #ALL} table without restating its settings. + */ + public boolean isWiderThan(MetaFieldsMode other) { + if (other == null) { + return false; + } + return (commitTimePopulated && !other.commitTimePopulated) + || (fileNamePopulated && !other.fileNamePopulated) + || (isRecordKeyPopulated() && !other.isRecordKeyPopulated()); + } +} diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java index 99f4cfdfd2893..db7ed17ad06d2 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java @@ -39,6 +39,7 @@ import org.apache.hudi.common.model.HoodieRecordPayload; import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.model.HoodieTimelineTimeZone; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.common.model.OverwriteNonDefaultsWithLatestAvroPayload; import org.apache.hudi.common.model.OverwriteWithLatestAvroPayload; import org.apache.hudi.common.model.PartialUpdateAvroPayload; @@ -327,12 +328,31 @@ public static final String getDefaultPayloadClassName() { .noDefaultValue() .withDocumentation("Base path of the dataset that needs to be bootstrapped as a Hudi table"); + /** + * @deprecated since 1.3.0, use {@link #META_FIELDS_MODE} instead. {@code true} maps to + * {@link MetaFieldsMode#ALL} and {@code false} maps to {@link MetaFieldsMode#NONE}. This property + * is still honored for tables written before {@code hoodie.meta.fields.mode} existed, but it is + * consulted only when the mode property is absent. + */ + @Deprecated public static final ConfigProperty POPULATE_META_FIELDS = ConfigProperty .key("hoodie.populate.meta.fields") .defaultValue(true) - .withDocumentation("When enabled, populates all meta fields. When disabled, no meta fields are populated " + .deprecatedAfter("1.2.0") + .withDocumentation("Deprecated — use hoodie.meta.fields.mode instead (true maps to ALL, false maps to NONE). " + + "When enabled, populates all meta fields. When disabled, no meta fields are populated " + "and incremental queries will not be functional. This is only meant to be used for append only/immutable data for batch processing"); + public static final ConfigProperty META_FIELDS_MODE = ConfigProperty + .key("hoodie.meta.fields.mode") + .defaultValue("") + .withDocumentation("Which Hudi meta columns are physically populated on disk. Allowed values are " + + "ALL, NONE, COMMIT_TIME_ONLY, FILE_NAME_ONLY and COMMIT_TIME_AND_FILE_NAME. This supersedes the " + + "deprecated hoodie.populate.meta.fields boolean, which is consulted only when this property is unset " + + "(true maps to ALL, false maps to NONE). Set only at table creation, " + + "via the hudi-cli, or during table upgrade — the property is immutable at runtime because it is a " + + "physical-storage decision baked into files at write time."); + public static final ConfigProperty KEY_GENERATOR_CLASS_NAME = ConfigProperty .key("hoodie.table.keygenerator.class") .noDefaultValue() @@ -1230,11 +1250,60 @@ public String getTimelinePath() { /** * @returns true is meta fields need to be populated. else returns false. + * + *

Derived from {@link #getMetaFieldsMode()} so that call sites still written against the + * deprecated boolean observe the same answer as the enum: only {@link MetaFieldsMode#ALL} + * populates every meta column. Selective modes report {@code false} here, which keeps + * key-dependent machinery (bloom filters, record-level index) correctly disabled. */ public boolean populateMetaFields() { + return getMetaFieldsMode().toLegacyPopulateMetaFields(); + } + + /** + * @return the raw, deprecated {@code hoodie.populate.meta.fields} value, used only as the + * fallback when {@link #META_FIELDS_MODE} is absent. Callers should use + * {@link #getMetaFieldsMode()} instead. + */ + private boolean legacyPopulateMetaFields() { return Boolean.parseBoolean(getStringOrDefault(POPULATE_META_FIELDS)); } + /** + * @return the {@link MetaFieldsMode} resolved from the on-disk properties. {@link #META_FIELDS_MODE} + * is the source of truth; tables written before that property existed fall back to + * {@link MetaFieldsMode#ALL} or {@link MetaFieldsMode#NONE} based on the deprecated + * {@link #POPULATE_META_FIELDS} boolean. + */ + public MetaFieldsMode getMetaFieldsMode() { + // Deliberately the two-argument form rather than resolve(this): the fallback must be the *raw* + // property, and populateMetaFields() on this class is itself derived from the mode. + return MetaFieldsMode.resolve(getStringOrDefault(META_FIELDS_MODE), legacyPopulateMetaFields()); + } + + /** + * @return true when the {@code _hoodie_commit_time} meta column is physically populated on disk. + */ + public boolean isCommitTimePopulated() { + return getMetaFieldsMode().isCommitTimePopulated(); + } + + /** + * @return true when the {@code _hoodie_file_name} meta column is physically populated on disk. + */ + public boolean isFileNamePopulated() { + return getMetaFieldsMode().isFileNamePopulated(); + } + + /** + * @return true when the {@code _hoodie_record_key} meta column is physically populated on disk. + * Only {@link MetaFieldsMode#ALL} populates the record-key column; every selective mode leaves + * it null. + */ + public boolean isRecordKeyPopulated() { + return getMetaFieldsMode().isRecordKeyPopulated(); + } + /** * @returns the record key field prop. */ diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableMetaClient.java b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableMetaClient.java index 0f3f8e7bf64b0..d5564ce9b2b68 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableMetaClient.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableMetaClient.java @@ -38,6 +38,7 @@ import org.apache.hudi.common.model.HoodieRecordPayload; import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.model.HoodieTimelineTimeZone; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.common.table.timeline.CommitMetadataSerDe; import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieArchivedTimeline; @@ -1013,6 +1014,7 @@ public static class TableBuilder { private String bootstrapBasePath; private Boolean bootstrapIndexEnable; private Boolean populateMetaFields; + private MetaFieldsMode metaFieldsMode; private String keyGeneratorClassProp; private String partitionValueExtractorClass; private String keyGeneratorType; @@ -1165,11 +1167,35 @@ public TableBuilder setBootstrapIndexEnable(Boolean bootstrapIndexEnable) { return this; } + /** + * @deprecated since 1.3.0, use {@link #setMetaFieldsMode(MetaFieldsMode)} instead + * ({@code true} maps to {@link MetaFieldsMode#ALL}, {@code false} to {@link MetaFieldsMode#NONE}). + */ + @Deprecated public TableBuilder setPopulateMetaFields(boolean populateMetaFields) { this.populateMetaFields = populateMetaFields; return this; } + public TableBuilder setMetaFieldsMode(MetaFieldsMode metaFieldsMode) { + this.metaFieldsMode = metaFieldsMode; + return this; + } + + /** + * Convenience overload that accepts the raw on-disk string (e.g. from properties files). + * Empty or null values leave the mode unset — the table then resolves to ALL or NONE from the + * deprecated populate.meta.fields boolean. + */ + public TableBuilder setMetaFieldsModeFromString(String rawMode) { + if (rawMode == null || rawMode.trim().isEmpty()) { + this.metaFieldsMode = null; + return this; + } + this.metaFieldsMode = MetaFieldsMode.parse(rawMode); + return this; + } + public TableBuilder setKeyGeneratorClassProp(String keyGeneratorClassProp) { this.keyGeneratorClassProp = keyGeneratorClassProp; return this; @@ -1384,6 +1410,9 @@ public TableBuilder fromProperties(Properties properties) { if (hoodieConfig.contains(HoodieTableConfig.POPULATE_META_FIELDS)) { setPopulateMetaFields(hoodieConfig.getBoolean(HoodieTableConfig.POPULATE_META_FIELDS)); } + if (hoodieConfig.contains(HoodieTableConfig.META_FIELDS_MODE)) { + setMetaFieldsModeFromString(hoodieConfig.getString(HoodieTableConfig.META_FIELDS_MODE)); + } if (hoodieConfig.contains(HoodieTableConfig.KEY_GENERATOR_CLASS_NAME)) { setKeyGeneratorClassProp(hoodieConfig.getString(HoodieTableConfig.KEY_GENERATOR_CLASS_NAME)); } else if (hoodieConfig.contains(HoodieTableConfig.KEY_GENERATOR_TYPE)) { @@ -1519,7 +1548,20 @@ public Properties build() { tableConfig.setValue(HoodieTableConfig.CDC_SUPPLEMENTAL_LOGGING_MODE, cdcSupplementalLoggingMode); } } - if (null != populateMetaFields) { + // hoodie.meta.fields.mode is the source of truth. When it is supplied, hoodie.properties must + // never contradict it: the legacy boolean is written from the mode (ALL -> true, every other + // mode -> false) rather than from whatever the caller passed. Otherwise a table written + // selectively could still record populate.meta.fields=true, and a pre-1.3.0 reader — which + // ignores the mode property entirely — would treat it as ALL. For NONE that is actively + // unsafe: an older incremental reader would be allowed to run against all-null commit times + // and silently return no rows. + if (null != metaFieldsMode) { + tableConfig.setValue(HoodieTableConfig.META_FIELDS_MODE, metaFieldsMode.name()); + tableConfig.setValue(HoodieTableConfig.POPULATE_META_FIELDS, + Boolean.toString(metaFieldsMode.toLegacyPopulateMetaFields())); + } else if (null != populateMetaFields) { + // No explicit mode: preserve pre-1.3.0 behavior and record only the legacy boolean, which + // resolves to ALL / NONE on read. tableConfig.setValue(HoodieTableConfig.POPULATE_META_FIELDS, Boolean.toString(populateMetaFields)); } if (null != keyGeneratorClassProp) { diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/TableSchemaResolver.java b/hudi-common/src/main/java/org/apache/hudi/common/table/TableSchemaResolver.java index 1e8988a742cbd..08c45e2be4e83 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/TableSchemaResolver.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/TableSchemaResolver.java @@ -23,6 +23,7 @@ import org.apache.hudi.common.model.HoodieCommitMetadata; import org.apache.hudi.common.model.HoodieLogFile; import org.apache.hudi.common.model.HoodieRecord; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.common.model.WriteOperationType; import org.apache.hudi.common.schema.HoodieSchema; import org.apache.hudi.common.schema.HoodieSchemaField; @@ -124,7 +125,11 @@ private Option getTableSchemaFromDataFileInternal() { * @throws Exception */ public HoodieSchema getTableSchema() throws Exception { - return getTableSchema(metaClient.getTableConfig().populateMetaFields()); + // Include meta fields whenever the table's meta-fields mode populates any of them. Under + // selective modes (COMMIT_TIME_ONLY / FILE_NAME_ONLY / COMMIT_TIME_AND_FILE_NAME) the meta + // columns exist as physical nullable Parquet columns even though populateMetaFields() is false, + // and read paths (e.g. incremental relations) must see them in the projected schema. + return getTableSchema(metaClient.getTableConfig().getMetaFieldsMode() != MetaFieldsMode.NONE); } /** @@ -148,7 +153,8 @@ public HoodieSchema getTableSchema(String timestamp) throws Exception { .filterCompletedInstants() .findInstantsBeforeOrEquals(timestamp) .lastInstant(); - return getTableSchemaInternal(metaClient.getTableConfig().populateMetaFields(), instant) + return getTableSchemaInternal( + metaClient.getTableConfig().getMetaFieldsMode() != MetaFieldsMode.NONE, instant) .orElseThrow(schemaNotFoundError()); } diff --git a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroFileWriterFactory.java b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroFileWriterFactory.java index 272de5d4a3103..db9558a01651b 100644 --- a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroFileWriterFactory.java +++ b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroFileWriterFactory.java @@ -69,7 +69,9 @@ public HoodieAvroFileWriterFactory(HoodieStorage storage) { protected HoodieFileWriter newParquetFileWriter( String instantTime, StoragePath path, HoodieConfig config, HoodieSchema schema, TaskContextSupplier taskContextSupplier) throws IOException { - boolean populateMetaFields = config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS); + org.apache.hudi.common.model.MetaFieldsMode metaFieldsMode = + org.apache.hudi.common.model.MetaFieldsMode.resolve(config); + boolean populateMetaFields = metaFieldsMode.toLegacyPopulateMetaFields(); Pair injectedConfigs = HoodieParquetConfigInjector.applyConfigInjector(path, storage.getConf(), config); StorageConfiguration storageConfiguration = injectedConfigs.getLeft(); @@ -89,7 +91,7 @@ protected HoodieFileWriter newParquetFileWriter( hoodieConfig.getLongOrDefault(HoodieStorageConfig.PARQUET_MAX_FILE_SIZE), storageConfiguration, hoodieConfig.getDoubleOrDefault(HoodieStorageConfig.PARQUET_COMPRESSION_RATIO_FRACTION), hoodieConfig.getBooleanOrDefault(HoodieStorageConfig.PARQUET_DICTIONARY_ENABLED)); - return new HoodieAvroParquetWriter(path, parquetConfig, instantTime, taskContextSupplier, populateMetaFields); + return new HoodieAvroParquetWriter(path, parquetConfig, instantTime, taskContextSupplier, metaFieldsMode); } protected HoodieFileWriter newParquetFileWriter( @@ -123,7 +125,13 @@ protected HoodieFileWriter newHFileFileWriter( HoodieAvroHFileReaderImplBase.KEY_FIELD_NAME, filter, config.getBoolean(HFILE_WRITER_TO_ALLOW_DUPLICATES)); - return new HoodieAvroHFileWriter(instantTime, path, hfileConfig, schema, taskContextSupplier, config.getBoolean(HoodieTableConfig.POPULATE_META_FIELDS)); + // Resolve through the mode like the parquet path above. HFile does not populate meta columns + // selectively, so a selective mode is treated as "record key not populated". Note getBoolean + // (unlike getBooleanOrDefault) returns null when neither property is set, which would NPE here. + boolean populateMetaFields = org.apache.hudi.common.model.MetaFieldsMode.resolve( + config.getStringOrDefault(HoodieTableConfig.META_FIELDS_MODE), + config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS)).isRecordKeyPopulated(); + return new HoodieAvroHFileWriter(instantTime, path, hfileConfig, schema, taskContextSupplier, populateMetaFields); } protected HoodieFileWriter newOrcFileWriter( diff --git a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroParquetWriter.java b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroParquetWriter.java index 1fb381b26b0ef..30a4e579deb75 100644 --- a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroParquetWriter.java +++ b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroParquetWriter.java @@ -23,10 +23,13 @@ import org.apache.hudi.common.config.HoodieParquetConfig; import org.apache.hudi.common.engine.TaskContextSupplier; import org.apache.hudi.common.model.HoodieKey; +import org.apache.hudi.common.model.HoodieRecord; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.core.io.storage.HoodieAvroFileWriter; import org.apache.hudi.io.hadoop.HoodieBaseParquetWriter; import org.apache.hudi.storage.StoragePath; +import org.apache.avro.generic.GenericRecord; import org.apache.avro.generic.IndexedRecord; import javax.annotation.concurrent.NotThreadSafe; @@ -48,39 +51,72 @@ public class HoodieAvroParquetWriter private final String fileName; private final String instantTime; private final TaskContextSupplier taskContextSupplier; - private final boolean populateMetaFields; + private final MetaFieldsMode metaFieldsMode; private final HoodieAvroWriteSupport writeSupport; + /** + * @deprecated since 1.3.0, use the {@link MetaFieldsMode} overload. Retained for existing callers + * that only distinguish all-or-nothing meta fields ({@code true} maps to {@link MetaFieldsMode#ALL}, + * {@code false} to {@link MetaFieldsMode#NONE}); it cannot express the selective modes. + */ + @Deprecated @SuppressWarnings({"unchecked", "rawtypes"}) public HoodieAvroParquetWriter(StoragePath file, HoodieParquetConfig parquetConfig, String instantTime, TaskContextSupplier taskContextSupplier, boolean populateMetaFields) throws IOException { + this(file, parquetConfig, instantTime, taskContextSupplier, + populateMetaFields ? MetaFieldsMode.ALL : MetaFieldsMode.NONE); + } + + @SuppressWarnings({"unchecked", "rawtypes"}) + public HoodieAvroParquetWriter(StoragePath file, + HoodieParquetConfig parquetConfig, + String instantTime, + TaskContextSupplier taskContextSupplier, + MetaFieldsMode metaFieldsMode) throws IOException { super(file, (HoodieParquetConfig) parquetConfig); this.fileName = file.getName(); this.writeSupport = parquetConfig.getWriteSupport(); this.instantTime = instantTime; this.taskContextSupplier = taskContextSupplier; - this.populateMetaFields = populateMetaFields; + this.metaFieldsMode = metaFieldsMode == null ? MetaFieldsMode.NONE : metaFieldsMode; } @Override public void writeAvroWithMetadata(HoodieKey key, IndexedRecord avroRecord) throws IOException { - if (populateMetaFields) { - prepRecordWithMetadata(key, avroRecord, instantTime, - taskContextSupplier.getPartitionIdSupplier().get(), getWrittenRecordCount(), fileName); - super.write(avroRecord); - writeSupport.add(key.getRecordKey()); - } else { - super.write(avroRecord); + switch (metaFieldsMode) { + case ALL: + prepRecordWithMetadata(key, avroRecord, instantTime, + taskContextSupplier.getPartitionIdSupplier().get(), getWrittenRecordCount(), fileName); + super.write(avroRecord); + writeSupport.add(key.getRecordKey()); + break; + case NONE: + super.write(avroRecord); + break; + default: + // Selective mode — populate only the opted-in columns. The other meta columns stay null, + // which Parquet stores as definition-level flags (zero data bytes). Bloom filter / + // record-key index population is intentionally skipped — that requires the record-key + // column, which is never populated in selective modes. + GenericRecord genericRecord = (GenericRecord) avroRecord; + if (metaFieldsMode.isCommitTimePopulated()) { + genericRecord.put(HoodieRecord.COMMIT_TIME_METADATA_FIELD, instantTime); + } + if (metaFieldsMode.isFileNamePopulated()) { + genericRecord.put(HoodieRecord.FILENAME_METADATA_FIELD, fileName); + } + super.write(avroRecord); + break; } } @Override public void writeAvro(String key, IndexedRecord object) throws IOException { super.write(object); - if (populateMetaFields) { + if (metaFieldsMode == MetaFieldsMode.ALL) { writeSupport.add(key); } } diff --git a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieMetaFieldsMode.java b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieMetaFieldsMode.java new file mode 100644 index 0000000000000..a7d884c808dac --- /dev/null +++ b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieMetaFieldsMode.java @@ -0,0 +1,174 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hudi.common.table; + +import org.apache.hudi.common.model.MetaFieldsMode; + +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Tests the meta-field-population modes exposed by {@link HoodieTableConfig} via the + * {@code hoodie.meta.fields.mode} property. The mode is materialized as a {@link MetaFieldsMode} + * enum resolved from both the legacy {@code hoodie.populate.meta.fields} boolean and the on-disk + * {@code hoodie.meta.fields.mode} property. + */ +class TestHoodieMetaFieldsMode { + + private static HoodieTableConfig configOf(Boolean populate, String mode) { + HoodieTableConfig cfg = new HoodieTableConfig(); + if (populate != null) { + cfg.setValue(HoodieTableConfig.POPULATE_META_FIELDS, String.valueOf(populate)); + } + if (mode != null) { + cfg.setValue(HoodieTableConfig.META_FIELDS_MODE, mode); + } + return cfg; + } + + @Test + void defaultsResolveToAllMode() { + HoodieTableConfig cfg = configOf(null, null); + assertTrue(cfg.populateMetaFields(), "populateMetaFields default must remain true"); + assertEquals(MetaFieldsMode.ALL, cfg.getMetaFieldsMode()); + assertTrue(cfg.isCommitTimePopulated()); + assertTrue(cfg.isFileNamePopulated()); + assertTrue(cfg.isRecordKeyPopulated()); + } + + @Test + void noneModeWhenPopulateFalseAndModeEmpty() { + HoodieTableConfig cfg = configOf(false, ""); + assertFalse(cfg.populateMetaFields()); + assertEquals(MetaFieldsMode.NONE, cfg.getMetaFieldsMode()); + assertFalse(cfg.isCommitTimePopulated()); + assertFalse(cfg.isFileNamePopulated()); + assertFalse(cfg.isRecordKeyPopulated()); + } + + @Test + void noneModeWhenPopulateFalseAndModeUnset() { + // Existing populate.meta.fields=false table without the mode property must resolve to NONE. + HoodieTableConfig cfg = configOf(false, null); + assertFalse(cfg.populateMetaFields()); + assertEquals(MetaFieldsMode.NONE, cfg.getMetaFieldsMode()); + assertFalse(cfg.isCommitTimePopulated()); + assertFalse(cfg.isFileNamePopulated()); + assertFalse(cfg.isRecordKeyPopulated()); + } + + @Test + void commitTimeOnlyMode() { + HoodieTableConfig cfg = configOf(false, MetaFieldsMode.COMMIT_TIME_ONLY.name()); + assertFalse(cfg.populateMetaFields()); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, cfg.getMetaFieldsMode()); + assertTrue(cfg.isCommitTimePopulated()); + assertFalse(cfg.isFileNamePopulated()); + assertFalse(cfg.isRecordKeyPopulated()); + } + + @Test + void fileNameOnlyMode() { + HoodieTableConfig cfg = configOf(false, MetaFieldsMode.FILE_NAME_ONLY.name()); + assertFalse(cfg.populateMetaFields()); + assertEquals(MetaFieldsMode.FILE_NAME_ONLY, cfg.getMetaFieldsMode()); + assertFalse(cfg.isCommitTimePopulated()); + assertTrue(cfg.isFileNamePopulated()); + assertFalse(cfg.isRecordKeyPopulated()); + } + + @Test + void commitTimeAndFileNameMode() { + HoodieTableConfig cfg = configOf(false, MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME.name()); + assertFalse(cfg.populateMetaFields()); + assertEquals(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME, cfg.getMetaFieldsMode()); + assertTrue(cfg.isCommitTimePopulated()); + assertTrue(cfg.isFileNamePopulated()); + assertFalse(cfg.isRecordKeyPopulated()); + } + + @Test + void modeWinsOverLegacyPopulateMetaFields() { + // hoodie.meta.fields.mode is the source of truth: when it is set, the deprecated + // populate.meta.fields boolean is not consulted, whichever way it points. + HoodieTableConfig cfg = configOf(true, MetaFieldsMode.COMMIT_TIME_ONLY.name()); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, cfg.getMetaFieldsMode()); + assertTrue(cfg.isCommitTimePopulated()); + assertFalse(cfg.isFileNamePopulated()); + // ...and populateMetaFields() is derived from the mode, so legacy call sites agree. + assertFalse(cfg.populateMetaFields()); + + HoodieTableConfig allWithLegacyFalse = configOf(false, MetaFieldsMode.ALL.name()); + assertEquals(MetaFieldsMode.ALL, allWithLegacyFalse.getMetaFieldsMode()); + assertTrue(allWithLegacyFalse.populateMetaFields()); + } + + @Test + void unknownTokenIsRejected() { + HoodieTableConfig cfg = configOf(false, "GARBAGE_VALUE"); + IllegalArgumentException ex = assertThrows(IllegalArgumentException.class, cfg::getMetaFieldsMode); + assertTrue(ex.getMessage().contains("GARBAGE_VALUE"), + "message must name the rejected value: " + ex.getMessage()); + assertTrue(ex.getMessage().contains("hoodie.meta.fields.mode"), + "message must name the property: " + ex.getMessage()); + } + + @Test + void isWiderThanRejectsEveryColumnAddingTransition() { + // Widening is what the write-client irreversibility guard must reject: any transition that + // adds a populated meta column, not just the all-or-nothing NONE -> ALL case. + assertTrue(MetaFieldsMode.COMMIT_TIME_ONLY.isWiderThan(MetaFieldsMode.NONE)); + assertTrue(MetaFieldsMode.FILE_NAME_ONLY.isWiderThan(MetaFieldsMode.NONE)); + assertTrue(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME.isWiderThan(MetaFieldsMode.COMMIT_TIME_ONLY)); + assertTrue(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME.isWiderThan(MetaFieldsMode.FILE_NAME_ONLY)); + assertTrue(MetaFieldsMode.ALL.isWiderThan(MetaFieldsMode.NONE)); + assertTrue(MetaFieldsMode.ALL.isWiderThan(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME)); + + // The two selective single-column modes each add a column the other lacks. + assertTrue(MetaFieldsMode.COMMIT_TIME_ONLY.isWiderThan(MetaFieldsMode.FILE_NAME_ONLY)); + assertTrue(MetaFieldsMode.FILE_NAME_ONLY.isWiderThan(MetaFieldsMode.COMMIT_TIME_ONLY)); + } + + @Test + void isWiderThanAllowsSameModeAndNarrowing() { + for (MetaFieldsMode mode : MetaFieldsMode.values()) { + assertFalse(mode.isWiderThan(mode), mode + " must not be wider than itself"); + } + // Narrowing drops columns from later commits, which is tolerated. + assertFalse(MetaFieldsMode.NONE.isWiderThan(MetaFieldsMode.ALL)); + assertFalse(MetaFieldsMode.COMMIT_TIME_ONLY.isWiderThan(MetaFieldsMode.ALL)); + assertFalse(MetaFieldsMode.COMMIT_TIME_ONLY.isWiderThan(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME)); + // A null on-disk mode carries no information, so nothing is "wider" than it. + assertFalse(MetaFieldsMode.ALL.isWiderThan(null)); + } + + @Test + void modeStringIsCaseInsensitiveAndTrimmed() { + // Whitespace around the value is tolerated, and casing does not have to match the enum. + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, configOf(false, " COMMIT_TIME_ONLY ").getMetaFieldsMode()); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, configOf(false, "commit_time_only").getMetaFieldsMode()); + assertEquals(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME, configOf(false, " Commit_Time_And_File_Name ").getMetaFieldsMode()); + HoodieTableConfig cfg = configOf(false, "CoMmIt_TiMe_OnLy"); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, cfg.getMetaFieldsMode()); + } +} diff --git a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableConfig.java b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableConfig.java index ff21ec822f821..273abd7a31e51 100644 --- a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableConfig.java +++ b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableConfig.java @@ -386,7 +386,7 @@ void testDropInvalidConfigs() { @Test void testDefinedTableConfigs() { List> configProperties = HoodieTableConfig.definedTableConfigs(); - assertEquals(45, configProperties.size()); + assertEquals(46, configProperties.size()); configProperties.forEach(c -> { assertNotNull(c); assertFalse(c.doc().isEmpty()); diff --git a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableMetaClient.java b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableMetaClient.java index 7e41c5578cc18..1c0ddbe091c69 100644 --- a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableMetaClient.java +++ b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/TestHoodieTableMetaClient.java @@ -22,6 +22,7 @@ import org.apache.hudi.common.model.HoodieIndexDefinition; import org.apache.hudi.common.model.HoodieIndexMetadata; import org.apache.hudi.common.model.HoodieTableType; +import org.apache.hudi.common.model.MetaFieldsMode; import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.timeline.HoodieTimeline; @@ -145,6 +146,54 @@ void testToString() throws IOException { assertNotEquals(metaClient1.toString(), new Object().toString()); } + @Test + void testMetaFieldsModeRewritesLegacyPopulateMetaFields() throws IOException { + // hoodie.properties must never record a legacy boolean that contradicts the mode: a pre-1.3.0 + // reader ignores hoodie.meta.fields.mode entirely and would otherwise treat a selectively + // written table as ALL. For NONE that is unsafe — an old incremental reader would run against + // all-null commit times and silently return no rows. + final String selectivePath = tempDir.toAbsolutePath() + Path.SEPARATOR + "mfm-selective"; + HoodieTableMetaClient selective = HoodieTableMetaClient.newTableBuilder() + .setTableType(HoodieTableType.COPY_ON_WRITE.name()) + .setTableName("mfm-selective") + .setPopulateMetaFields(true) + .setMetaFieldsMode(MetaFieldsMode.COMMIT_TIME_ONLY) + .initTable(this.metaClient.getStorageConf(), selectivePath); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, selective.getTableConfig().getMetaFieldsMode()); + assertFalse(selective.getTableConfig().populateMetaFields(), + "caller-supplied populate.meta.fields=true must be overridden by the mode"); + + final String nonePath = tempDir.toAbsolutePath() + Path.SEPARATOR + "mfm-none"; + HoodieTableMetaClient none = HoodieTableMetaClient.newTableBuilder() + .setTableType(HoodieTableType.COPY_ON_WRITE.name()) + .setTableName("mfm-none") + .setPopulateMetaFields(true) + .setMetaFieldsMode(MetaFieldsMode.NONE) + .initTable(this.metaClient.getStorageConf(), nonePath); + assertEquals(MetaFieldsMode.NONE, none.getTableConfig().getMetaFieldsMode()); + assertFalse(none.getTableConfig().populateMetaFields()); + + final String allPath = tempDir.toAbsolutePath() + Path.SEPARATOR + "mfm-all"; + HoodieTableMetaClient all = HoodieTableMetaClient.newTableBuilder() + .setTableType(HoodieTableType.COPY_ON_WRITE.name()) + .setTableName("mfm-all") + .setPopulateMetaFields(false) + .setMetaFieldsMode(MetaFieldsMode.ALL) + .initTable(this.metaClient.getStorageConf(), allPath); + assertEquals(MetaFieldsMode.ALL, all.getTableConfig().getMetaFieldsMode()); + assertTrue(all.getTableConfig().populateMetaFields()); + + // No explicit mode: pre-1.3.0 behavior preserved, only the legacy boolean is recorded. + final String legacyPath = tempDir.toAbsolutePath() + Path.SEPARATOR + "mfm-legacy"; + HoodieTableMetaClient legacy = HoodieTableMetaClient.newTableBuilder() + .setTableType(HoodieTableType.COPY_ON_WRITE.name()) + .setTableName("mfm-legacy") + .setPopulateMetaFields(false) + .initTable(this.metaClient.getStorageConf(), legacyPath); + assertEquals(MetaFieldsMode.NONE, legacy.getTableConfig().getMetaFieldsMode()); + assertFalse(legacy.getTableConfig().populateMetaFields()); + } + @Test void testTableVersion() throws IOException { final String basePath = tempDir.toAbsolutePath() + Path.SEPARATOR + "t1"; diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/commit/BaseDatasetBulkInsertCommitActionExecutor.java b/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/commit/BaseDatasetBulkInsertCommitActionExecutor.java index 8b6466d5117d6..de1505037a6d8 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/commit/BaseDatasetBulkInsertCommitActionExecutor.java +++ b/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/commit/BaseDatasetBulkInsertCommitActionExecutor.java @@ -30,7 +30,6 @@ import org.apache.hudi.common.model.HoodieFileGroupId; import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.WriteOperationType; -import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.util.CommitUtils; import org.apache.hudi.common.util.Option; @@ -114,7 +113,10 @@ public final HoodieWriteResult execute(Dataset records, boolean isTablePart throw new HoodieException("Dropping duplicates with bulk_insert in row writer path is not supported yet"); } - boolean populateMetaFields = writeConfig.getBoolean(HoodieTableConfig.POPULATE_META_FIELDS); + // Resolve through the mode rather than reading the deprecated boolean directly: the sort / + // bucket partitioners below key off _hoodie_record_key and _hoodie_partition_path, which are + // null under every selective mode. + boolean populateMetaFields = writeConfig.getMetaFieldsMode().isRecordKeyPopulated(); preExecute(); BulkInsertPartitioner> bulkInsertPartitionerRows = getPartitioner(populateMetaFields, isTablePartitioned); diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/AutoRecordKeyGenerationUtils.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/AutoRecordKeyGenerationUtils.scala index 5d8bc1ee405ee..d32d242404571 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/AutoRecordKeyGenerationUtils.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/AutoRecordKeyGenerationUtils.scala @@ -21,6 +21,7 @@ package org.apache.hudi import org.apache.hudi.DataSourceWriteOptions.INSERT_DROP_DUPS import org.apache.hudi.common.config.HoodieConfig +import org.apache.hudi.common.model.MetaFieldsMode import org.apache.hudi.common.table.HoodieTableConfig import org.apache.hudi.common.util.{ConfigUtils, StringUtils} import org.apache.hudi.config.HoodieWriteConfig @@ -43,9 +44,20 @@ object AutoRecordKeyGenerationUtils { if (hoodieConfig.getBoolean(INSERT_DROP_DUPS)) { throw new HoodieKeyGeneratorException("Enabling " + INSERT_DROP_DUPS.key() + " is not supported with auto generation of record keys ") } - // virtual keys are not supported with auto generation of record keys. - if (!parameters.getOrElse(HoodieTableConfig.POPULATE_META_FIELDS.key(), HoodieTableConfig.POPULATE_META_FIELDS.defaultValue().toString).toBoolean) { - throw new HoodieKeyGeneratorException("Disabling " + HoodieTableConfig.POPULATE_META_FIELDS.key() + " is not supported with auto generation of record keys") + // virtual keys are not supported with auto generation of record keys. Resolve the mode rather + // than the deprecated boolean alone — a selective mode also leaves _hoodie_record_key + // unpopulated, so auto-generated keys would be computed and then discarded. + val metaFieldsMode = MetaFieldsMode.resolve( + parameters.getOrElse(HoodieTableConfig.META_FIELDS_MODE.key(), null), + parameters.getOrElse(HoodieTableConfig.POPULATE_META_FIELDS.key(), + HoodieTableConfig.POPULATE_META_FIELDS.defaultValue().toString).toBoolean) + if (!metaFieldsMode.isRecordKeyPopulated) { + // Name whichever property the user actually set, so the error points at the config to change. + val offendingKey = + if (parameters.contains(HoodieTableConfig.META_FIELDS_MODE.key())) HoodieTableConfig.META_FIELDS_MODE.key() + else HoodieTableConfig.POPULATE_META_FIELDS.key() + throw new HoodieKeyGeneratorException(offendingKey + " is not supported with auto generation of record keys" + + " (resolved meta fields mode " + metaFieldsMode + " does not populate _hoodie_record_key)") } val orderingFieldsStr = ConfigUtils.getOrderingFieldsStrDuringWrite(hoodieConfig.getProps) if (StringUtils.nonEmpty(orderingFieldsStr)) { diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/DataSourceOptions.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/DataSourceOptions.scala index b2548029e6216..d84492e250982 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/DataSourceOptions.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/DataSourceOptions.scala @@ -21,7 +21,7 @@ import org.apache.hudi.DataSourceReadOptions.{QUERY_TYPE, QUERY_TYPE_READ_OPTIMI import org.apache.hudi.HoodieConversionUtils.toScalaOption import org.apache.hudi.common.config._ import org.apache.hudi.common.fs.ConsistencyGuardConfig -import org.apache.hudi.common.model.{HoodieTableType, WriteOperationType} +import org.apache.hudi.common.model.{HoodieTableType, MetaFieldsMode, WriteOperationType} import org.apache.hudi.common.table.HoodieTableConfig import org.apache.hudi.common.util.{Option, StringUtils} import org.apache.hudi.common.util.ConfigUtils.{DELTA_STREAMER_CONFIG_PREFIX, IS_QUERY_AS_RO_TABLE, STREAMER_CONFIG_PREFIX} @@ -548,8 +548,12 @@ object DataSourceWriteOptions { .defaultValue("true") .withInferFunction( JFunction.toJavaFunction((config: HoodieConfig) => { + // Resolve the meta-fields mode rather than reading the deprecated boolean: a selective mode + // leaves _hoodie_record_key unpopulated just as populate.meta.fields=false does, and must + // disable the row writer for the same reason. if (config.getString(OPERATION) == WriteOperationType.BULK_INSERT.value - && !config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS) + && !MetaFieldsMode.resolve(config.getStringOrDefault(HoodieTableConfig.META_FIELDS_MODE), + config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS)).isRecordKeyPopulated && config.getBooleanOrDefault(HoodieWriteConfig.COMBINE_BEFORE_INSERT)) { // need to turn off row writing for BULK_INSERT without meta fields with turned on COMBINE_BEFORE_INSERT to prevent shortcutting and ignoring COMBINE_BEFORE_INSERT setting Option.of("false") diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieSparkSqlWriter.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieSparkSqlWriter.scala index 37f8e569a0e7a..5d66069733589 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieSparkSqlWriter.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieSparkSqlWriter.scala @@ -284,6 +284,7 @@ class HoodieSparkSqlWriterInternal { val baseFileFormat = hoodieConfig.getStringOrDefault(HoodieTableConfig.BASE_FILE_FORMAT) val archiveLogFolder = hoodieConfig.getStringOrDefault(HoodieTableConfig.TIMELINE_HISTORY_PATH) val populateMetaFields = hoodieConfig.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS) + val metaFieldsMode = hoodieConfig.getStringOrDefault(HoodieTableConfig.META_FIELDS_MODE) val useBaseFormatMetaFile = hoodieConfig.getBooleanOrDefault(HoodieTableConfig.PARTITION_METAFILE_USE_BASE_FORMAT); val payloadClass = hoodieConfig.getString(DataSourceWriteOptions.PAYLOAD_CLASS_NAME) val recordMergeStrategyId = hoodieConfig.getString(DataSourceWriteOptions.RECORD_MERGE_STRATEGY_ID) @@ -308,6 +309,7 @@ class HoodieSparkSqlWriterInternal { .setOrderingFields(ConfigUtils.getOrderingFieldsStrDuringWrite(optParams.asJava)) .setPartitionFields(partitionColumnsForKeyGenerator) .setPopulateMetaFields(populateMetaFields) + .setMetaFieldsModeFromString(metaFieldsMode) .setRecordKeyFields(hoodieConfig.getString(RECORDKEY_FIELD)) .setSecondaryKeyFields(hoodieConfig.getString(SECONDARYKEY_COLUMN_NAME)) .setCDCEnabled(hoodieConfig.getBooleanOrDefault(HoodieTableConfig.CDC_ENABLED)) @@ -754,6 +756,10 @@ class HoodieSparkSqlWriterInternal { HoodieTableConfig.POPULATE_META_FIELDS.key(), String.valueOf(HoodieTableConfig.POPULATE_META_FIELDS.defaultValue()) )) + val metaFieldsMode = parameters.getOrElse( + HoodieTableConfig.META_FIELDS_MODE.key(), + HoodieTableConfig.META_FIELDS_MODE.defaultValue() + ) val baseFileFormat = hoodieConfig.getStringOrDefault(HoodieTableConfig.BASE_FILE_FORMAT) val useBaseFormatMetaFile = java.lang.Boolean.parseBoolean(parameters.getOrElse( HoodieTableConfig.PARTITION_METAFILE_USE_BASE_FORMAT.key(), @@ -780,6 +786,7 @@ class HoodieSparkSqlWriterInternal { .setCDCEnabled(hoodieConfig.getBooleanOrDefault(HoodieTableConfig.CDC_ENABLED)) .setCDCSupplementalLoggingMode(hoodieConfig.getStringOrDefault(HoodieTableConfig.CDC_SUPPLEMENTAL_LOGGING_MODE)) .setPopulateMetaFields(populateMetaFields) + .setMetaFieldsModeFromString(metaFieldsMode) .setKeyGeneratorClassProp(keyGenProp) .setPartitionValueExtractorClass(partitionValueExtractorClassName) .set(timestampKeyGeneratorConfigs.asJava.asInstanceOf[java.util.Map[String, Object]]) diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieWriterUtils.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieWriterUtils.scala index 54459ff61e00a..e0cbe5858849d 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieWriterUtils.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieWriterUtils.scala @@ -352,6 +352,23 @@ object HoodieWriterUtils { diffConfigs.append(s"${HoodieTableConfig.RECORD_MERGE_STRATEGY_ID}:\t$mergeStrategyId\tnull\n") } } + + // hoodie.meta.fields.mode is a physical-storage decision baked into files at write time. + // Changing it at runtime would silently produce mixed-mode files whose incremental / file + // pruning behavior differs between old and new commits. The default loop above only flags + // the mismatch when the on-disk value is non-null, so an older table with the property + // absent from hoodie.properties would let a null → non-empty transition slip through + // (silent-drop risk on pre-enablement commits). Guard the null → selective-mode case + // explicitly here. Set it only at table creation, via the hudi-cli, or during table upgrade; + // otherwise the only way to change it is to recreate the table. + val paramsMetaFieldsMode = params.getOrElse(HoodieTableConfig.META_FIELDS_MODE.key(), "") + val onDiskMetaFieldsMode = tableConfig.getString(HoodieTableConfig.META_FIELDS_MODE) + if (paramsMetaFieldsMode.nonEmpty && (onDiskMetaFieldsMode == null || onDiskMetaFieldsMode.isEmpty)) { + diffConfigs.append( + s"${HoodieTableConfig.META_FIELDS_MODE.key()}:\t$paramsMetaFieldsMode\t${if (onDiskMetaFieldsMode == null) "null" else "\"\""}" + + " (immutable at runtime; set only at table creation / hudi-cli / upgrade; " + + "existing tables must be recreated to change this)\n") + } } if (diffConfigs.nonEmpty) { diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/IncrementalRelationV1.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/IncrementalRelationV1.scala index 50db4c678c8f8..31c9620a7b0b6 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/IncrementalRelationV1.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/IncrementalRelationV1.scala @@ -89,8 +89,12 @@ class IncrementalRelationV1(val sqlContext: SQLContext, s"option ${DataSourceReadOptions.START_COMMIT.key}") } - if (!metaClient.getTableConfig.populateMetaFields()) { - throw new HoodieException("Incremental queries are not supported when meta fields are disabled") + if (!metaClient.getTableConfig.isCommitTimePopulated()) { + throw new HoodieException("Incremental queries are not supported when _hoodie_commit_time is not populated. " + + "hoodie.meta.fields.mode is a physical-storage decision baked into files at write time and cannot be " + + "changed by flipping write options — setting it only takes effect at table creation. To enable incremental " + + "queries on this table, recreate it with hoodie.populate.meta.fields=true or hoodie.meta.fields.mode=COMMIT_TIME_ONLY " + + "(or COMMIT_TIME_AND_FILE_NAME).") } private val useEndInstantSchema = optParams.getOrElse(INCREMENTAL_READ_SCHEMA_USE_END_INSTANTTIME.key, diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/IncrementalRelationV2.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/IncrementalRelationV2.scala index 308ca87a47592..47fbd3799ffee 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/IncrementalRelationV2.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/IncrementalRelationV2.scala @@ -78,8 +78,12 @@ class IncrementalRelationV2(val sqlContext: SQLContext, s"option ${DataSourceReadOptions.START_COMMIT.key}") } - if (!metaClient.getTableConfig.populateMetaFields()) { - throw new HoodieException("Incremental queries are not supported when meta fields are disabled") + if (!metaClient.getTableConfig.isCommitTimePopulated()) { + throw new HoodieException("Incremental queries are not supported when _hoodie_commit_time is not populated. " + + "hoodie.meta.fields.mode is a physical-storage decision baked into files at write time and cannot be " + + "changed by flipping write options — setting it only takes effect at table creation. To enable incremental " + + "queries on this table, recreate it with hoodie.populate.meta.fields=true or hoodie.meta.fields.mode=COMMIT_TIME_ONLY " + + "(or COMMIT_TIME_AND_FILE_NAME).") } private val queryContext: IncrementalQueryAnalyzer.QueryContext = diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/MergeOnReadIncrementalRelationV1.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/MergeOnReadIncrementalRelationV1.scala index e0fc7ceb6824b..f43071ca2ee89 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/MergeOnReadIncrementalRelationV1.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/MergeOnReadIncrementalRelationV1.scala @@ -267,8 +267,14 @@ trait HoodieIncrementalRelationV1Trait extends HoodieBaseRelation { s"option ${DataSourceReadOptions.START_COMMIT.key}") } + // MoR incremental relies on _hoodie_commit_time being present in BOTH base files AND log + // records. The base-file writer respects hoodie.meta.fields.mode, but the log-write path + // (HoodieAppendHandle) does not yet — until that gap is closed, MoR incremental must require + // populate.meta.fields=true to avoid silently dropping log-file rows whose commit_time is null. if (!this.tableConfig.populateMetaFields()) { - throw new HoodieException("Incremental queries are not supported when meta fields are disabled") + throw new HoodieException("Incremental queries on MoR tables are not supported when " + + "hoodie.populate.meta.fields=false. Selective meta-field modes (hoodie.meta.fields.mode) " + + "are supported for CoW only in this release; MoR support is tracked as a follow-up.") } if (hollowCommitHandling == USE_TRANSITION_TIME && fullTableScan) { diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/MergeOnReadIncrementalRelationV2.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/MergeOnReadIncrementalRelationV2.scala index aea594d9157ca..5e9a611bc6637 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/MergeOnReadIncrementalRelationV2.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/MergeOnReadIncrementalRelationV2.scala @@ -257,8 +257,14 @@ trait HoodieIncrementalRelationV2Trait extends HoodieBaseRelation { s"option ${DataSourceReadOptions.START_COMMIT.key}") } + // MoR incremental relies on _hoodie_commit_time being present in BOTH base files AND log + // records. The base-file writer respects hoodie.meta.fields.mode, but the log-write path + // (HoodieAppendHandle) does not yet — until that gap is closed, MoR incremental must require + // populate.meta.fields=true to avoid silently dropping log-file rows whose commit_time is null. if (!this.tableConfig.populateMetaFields()) { - throw new HoodieException("Incremental queries are not supported when meta fields are disabled") + throw new HoodieException("Incremental queries on MoR tables are not supported when " + + "hoodie.populate.meta.fields=false. Selective meta-field modes (hoodie.meta.fields.mode) " + + "are supported for CoW only in this release; MoR support is tracked as a follow-up.") } } diff --git a/hudi-spark-datasource/hudi-spark/src/main/java/org/apache/hudi/cli/BootstrapExecutorUtils.java b/hudi-spark-datasource/hudi-spark/src/main/java/org/apache/hudi/cli/BootstrapExecutorUtils.java index b4d6c47e87011..462e66a4aa95e 100644 --- a/hudi-spark-datasource/hudi-spark/src/main/java/org/apache/hudi/cli/BootstrapExecutorUtils.java +++ b/hudi-spark-datasource/hudi-spark/src/main/java/org/apache/hudi/cli/BootstrapExecutorUtils.java @@ -238,6 +238,8 @@ private void initializeTable() throws IOException { .setOrderingFields(ConfigUtils.getOrderingFieldsStrDuringWrite(props)) .setPopulateMetaFields(props.getBoolean( POPULATE_META_FIELDS.key(), POPULATE_META_FIELDS.defaultValue())) + .setMetaFieldsModeFromString(props.getString( + HoodieTableConfig.META_FIELDS_MODE.key(), HoodieTableConfig.META_FIELDS_MODE.defaultValue())) .setArchiveLogFolder(props.getString( TIMELINE_HISTORY_PATH.key(), TIMELINE_HISTORY_PATH.defaultValue())) .setPayloadClassName(cfg.payloadClass) diff --git a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestHoodieDatasetBulkInsertHelper.java b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestHoodieDatasetBulkInsertHelper.java index 946fb61c3b3a8..41cabe3d886df 100644 --- a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestHoodieDatasetBulkInsertHelper.java +++ b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestHoodieDatasetBulkInsertHelper.java @@ -188,11 +188,13 @@ public void testBulkInsertHelperNoMetaFields() { } result.toJavaRDD().foreach(entry -> { - assertTrue(entry.get(resultSchema.fieldIndex(HoodieRecord.RECORD_KEY_METADATA_FIELD)).equals("")); - assertTrue(entry.get(resultSchema.fieldIndex(HoodieRecord.PARTITION_PATH_METADATA_FIELD)).equals("")); - assertTrue(entry.get(resultSchema.fieldIndex(HoodieRecord.COMMIT_SEQNO_METADATA_FIELD)).equals("")); - assertTrue(entry.get(resultSchema.fieldIndex(HoodieRecord.COMMIT_TIME_METADATA_FIELD)).equals("")); - assertTrue(entry.get(resultSchema.fieldIndex(HoodieRecord.FILENAME_METADATA_FIELD)).equals("")); + // The stub meta columns are now written as null (nullable=true) rather than empty strings — + // see HoodieDatasetBulkInsertHelper. Handle handle-level population fills them in later. + assertTrue(entry.isNullAt(resultSchema.fieldIndex(HoodieRecord.RECORD_KEY_METADATA_FIELD))); + assertTrue(entry.isNullAt(resultSchema.fieldIndex(HoodieRecord.PARTITION_PATH_METADATA_FIELD))); + assertTrue(entry.isNullAt(resultSchema.fieldIndex(HoodieRecord.COMMIT_SEQNO_METADATA_FIELD))); + assertTrue(entry.isNullAt(resultSchema.fieldIndex(HoodieRecord.COMMIT_TIME_METADATA_FIELD))); + assertTrue(entry.isNullAt(resultSchema.fieldIndex(HoodieRecord.FILENAME_METADATA_FIELD))); }); Dataset trimmedOutput = result.drop(HoodieRecord.PARTITION_PATH_METADATA_FIELD).drop(HoodieRecord.RECORD_KEY_METADATA_FIELD) diff --git a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestMetaFieldsMode.java b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestMetaFieldsMode.java new file mode 100644 index 0000000000000..f1792f6366ba7 --- /dev/null +++ b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestMetaFieldsMode.java @@ -0,0 +1,469 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hudi.functional; + +import org.apache.hudi.DataSourceWriteOptions; +import org.apache.hudi.SparkAdapterSupport$; +import org.apache.hudi.common.config.HoodieMetadataConfig; +import org.apache.hudi.common.model.HoodieRecord; +import org.apache.hudi.common.model.MetaFieldsMode; +import org.apache.hudi.common.table.HoodieTableConfig; +import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.testutils.SparkClientFunctionalTestHarness; + +import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Row; +import org.apache.spark.sql.RowFactory; +import org.apache.spark.sql.SaveMode; +import org.apache.spark.sql.types.DataTypes; +import org.apache.spark.sql.types.StructField; +import org.apache.spark.sql.types.StructType; +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Spark-datasource end-to-end tests for the {@code hoodie.meta.fields.mode} property on CoW tables. + * Every {@link MetaFieldsMode} value is exercised via a write / re-read round trip; on-disk column + * population is verified by reading the parquet files back and inspecting the meta-column values. + */ +class TestMetaFieldsMode extends SparkClientFunctionalTestHarness { + + private static StructType simpleSchema() { + return DataTypes.createStructType(new StructField[]{ + DataTypes.createStructField("column1", DataTypes.StringType, true), + DataTypes.createStructField("column2", DataTypes.StringType, true), + DataTypes.createStructField("column3", DataTypes.StringType, true) + }).asNullable(); + } + + private Map baseOptions() { + Map opts = new HashMap<>(); + opts.put(DataSourceWriteOptions.RECORDKEY_FIELD().key(), "column1"); + opts.put(DataSourceWriteOptions.PARTITIONPATH_FIELD().key(), "column2"); + opts.put(DataSourceWriteOptions.ORDERING_FIELDS().key(), "column3"); + opts.put(HoodieTableConfig.NAME.key(), "test_meta_fields_mode"); + opts.put(DataSourceWriteOptions.TABLE_TYPE().key(), "COPY_ON_WRITE"); + opts.put(HoodieMetadataConfig.ENABLE.key(), "false"); + return opts; + } + + private void writeRows(List records, StructType schema, Map options, String path, SaveMode mode) { + spark().createDataset(records, + SparkAdapterSupport$.MODULE$.sparkAdapter().getCatalystExpressionUtils().getEncoder(schema)) + .write() + .format("hudi") + .options(options) + .mode(mode) + .save(path); + } + + private HoodieTableConfig writeSampleAndGetTableConfig(Map options, String path) { + writeRows(Arrays.asList( + RowFactory.create("k1", "p1", "v1"), + RowFactory.create("k2", "p1", "v2")), + simpleSchema(), options, path, SaveMode.Overwrite); + HoodieTableMetaClient metaClient = + HoodieTableMetaClient.builder().setBasePath(path).setConf(storageConf()).build(); + return metaClient.getTableConfig(); + } + + /** + * End-to-end assertion of the on-disk meta columns after a write. Reads the parquet files back + * (bypassing Hudi's own read path so we see the raw column values) and asserts which meta + * columns are non-null. + */ + private void assertMetaColumnPopulation(String path, MetaFieldsMode expectedMode) { + Dataset raw = spark().read().parquet(path + "/*/*.parquet"); + Row first = raw.select( + HoodieRecord.COMMIT_TIME_METADATA_FIELD, + HoodieRecord.COMMIT_SEQNO_METADATA_FIELD, + HoodieRecord.RECORD_KEY_METADATA_FIELD, + HoodieRecord.PARTITION_PATH_METADATA_FIELD, + HoodieRecord.FILENAME_METADATA_FIELD).first(); + + if (expectedMode.isCommitTimePopulated()) { + assertNotNull(first.get(0), "expected _hoodie_commit_time to be populated for mode " + expectedMode); + } else { + assertNull(first.get(0), "expected _hoodie_commit_time to be null for mode " + expectedMode); + } + if (expectedMode.isFileNamePopulated()) { + assertNotNull(first.get(4), "expected _hoodie_file_name to be populated for mode " + expectedMode); + } else { + assertNull(first.get(4), "expected _hoodie_file_name to be null for mode " + expectedMode); + } + // Record key, partition path, and commit seq no are ALL-only. + if (expectedMode == MetaFieldsMode.ALL) { + assertNotNull(first.get(2), "record key must be populated in ALL mode"); + assertNotNull(first.get(3), "partition path must be populated in ALL mode"); + assertNotNull(first.get(1), "commit seq no must be populated in ALL mode"); + } else { + assertNull(first.get(2), "record key must be null outside ALL mode, got: " + first.get(2)); + assertNull(first.get(3), "partition path must be null outside ALL mode, got: " + first.get(3)); + assertNull(first.get(1), "commit seq no must be null outside ALL mode, got: " + first.get(1)); + } + } + + @Test + void allModePersistsAndPopulatesAllColumns() { + Map options = baseOptions(); + // ALL is the default; no need to set the mode explicitly. + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + + assertTrue(tc.populateMetaFields()); + assertEquals(MetaFieldsMode.ALL, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.ALL); + } + + @Test + void noneModePersistsAndLeavesAllColumnsNull() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + + assertFalse(tc.populateMetaFields()); + assertEquals(MetaFieldsMode.NONE, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.NONE); + } + + @Test + void commitTimeOnlyModePopulatesOnlyCommitTime() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.COMMIT_TIME_ONLY.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY.name(), + tc.getProps().getProperty(HoodieTableConfig.META_FIELDS_MODE.key())); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.COMMIT_TIME_ONLY); + } + + @Test + void fileNameOnlyModePopulatesOnlyFileName() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.FILE_NAME_ONLY.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + + assertEquals(MetaFieldsMode.FILE_NAME_ONLY, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.FILE_NAME_ONLY); + } + + @Test + void commitTimeAndFileNameModePopulatesBoth() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + + assertEquals(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME); + } + + @Test + void selectiveModeWinsOverLegacyPopulateTrue() { + // hoodie.meta.fields.mode is the source of truth: an explicit mode is honored regardless of + // the deprecated boolean, so this combination is no longer ambiguous and is not rejected. + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "true"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.COMMIT_TIME_ONLY.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.COMMIT_TIME_ONLY); + // ...and hoodie.properties must not contradict the mode. A pre-1.3.0 reader ignores the mode + // property entirely, so leaving populate.meta.fields=true here would make it treat a + // selectively-written table as ALL. + assertFalse(tc.populateMetaFields(), + "legacy populate.meta.fields must be derived from the mode, not carried through verbatim"); + } + + @Test + void noneModePersistsLegacyBooleanAsFalse() { + // The unsafe case: an old incremental reader that sees populate.meta.fields=true on a NONE + // table would run against all-null commit times and silently return zero rows. + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "true"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.NONE.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + + assertEquals(MetaFieldsMode.NONE, tc.getMetaFieldsMode()); + assertFalse(tc.populateMetaFields(), + "NONE must persist populate.meta.fields=false so pre-1.3.0 readers do not treat it as ALL"); + } + + @Test + void allModePersistsLegacyBooleanAsTrue() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.ALL.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + + assertEquals(MetaFieldsMode.ALL, tc.getMetaFieldsMode()); + assertTrue(tc.populateMetaFields(), + "ALL must persist populate.meta.fields=true for pre-1.3.0 readers"); + } + + @Test + void unknownModeValueIsRejected() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), "SOMETHING_BOGUS"); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + Throwable thrown = assertThrows(Throwable.class, () -> + writeRows(Collections.singletonList(RowFactory.create("k1", "p1", "v1")), + simpleSchema(), options, basePath(), SaveMode.Overwrite)); + + String rootMessage = rootMessageOf(thrown); + assertTrue(rootMessage.contains("SOMETHING_BOGUS"), + "Expected error to name the rejected value, got: " + rootMessage); + } + + // ------------------------------------------------------------------------- + // Non-row-writer path coverage. Bulk insert with row.writer.enable=false forces the + // HoodieAvroParquetWriter path (via HoodieCreateHandle) instead of the internal-row writer path. + // Both paths must respect the mode identically. + // ------------------------------------------------------------------------- + + @Test + void nonRowWriterPathAllMode() { + Map options = baseOptions(); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.INSERT_OPERATION_OPT_VAL()); + options.put("hoodie.datasource.write.row.writer.enable", "false"); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + assertEquals(MetaFieldsMode.ALL, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.ALL); + } + + @Test + void nonRowWriterPathNoneMode() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.INSERT_OPERATION_OPT_VAL()); + options.put("hoodie.datasource.write.row.writer.enable", "false"); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + assertEquals(MetaFieldsMode.NONE, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.NONE); + } + + @Test + void nonRowWriterPathCommitTimeOnly() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.COMMIT_TIME_ONLY.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.INSERT_OPERATION_OPT_VAL()); + options.put("hoodie.datasource.write.row.writer.enable", "false"); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.COMMIT_TIME_ONLY); + } + + @Test + void nonRowWriterPathFileNameOnly() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.FILE_NAME_ONLY.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.INSERT_OPERATION_OPT_VAL()); + options.put("hoodie.datasource.write.row.writer.enable", "false"); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + assertEquals(MetaFieldsMode.FILE_NAME_ONLY, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.FILE_NAME_ONLY); + } + + @Test + void nonRowWriterPathCommitTimeAndFileName() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.INSERT_OPERATION_OPT_VAL()); + options.put("hoodie.datasource.write.row.writer.enable", "false"); + + HoodieTableConfig tc = writeSampleAndGetTableConfig(options, basePath()); + assertEquals(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME, tc.getMetaFieldsMode()); + assertMetaColumnPopulation(basePath(), MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME); + } + + // ------------------------------------------------------------------------- + // Clustering coverage. Inline clustering rewrites files through the create/merge handles which + // delegate to the same underlying HoodieAvroParquetWriter / HoodieRowCreateHandle we exercise + // in the write tests. Verifies clustered files preserve the mode's column population semantics. + // ------------------------------------------------------------------------- + + @Test + void clusteringPreservesCommitTimeOnlyMode() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.COMMIT_TIME_ONLY.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + // Inline clustering after each write. + options.put("hoodie.clustering.inline", "true"); + options.put("hoodie.clustering.inline.max.commits", "1"); + options.put("hoodie.clustering.plan.strategy.target.file.max.bytes", "10485760"); + options.put("hoodie.clustering.plan.strategy.small.file.limit", "10485760"); + + writeRows(Arrays.asList( + RowFactory.create("k1", "p1", "v1"), + RowFactory.create("k2", "p1", "v2"), + RowFactory.create("k3", "p1", "v3")), + simpleSchema(), options, basePath(), SaveMode.Overwrite); + + HoodieTableMetaClient metaClient = + HoodieTableMetaClient.builder().setBasePath(basePath()).setConf(storageConf()).build(); + assertEquals(MetaFieldsMode.COMMIT_TIME_ONLY, metaClient.getTableConfig().getMetaFieldsMode()); + // After clustering, files are rewritten — verify the rewritten files still respect the mode. + assertMetaColumnPopulation(basePath(), MetaFieldsMode.COMMIT_TIME_ONLY); + } + + @Test + void clusteringPreservesFileNameOnlyMode() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.FILE_NAME_ONLY.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + options.put("hoodie.clustering.inline", "true"); + options.put("hoodie.clustering.inline.max.commits", "1"); + options.put("hoodie.clustering.plan.strategy.target.file.max.bytes", "10485760"); + options.put("hoodie.clustering.plan.strategy.small.file.limit", "10485760"); + + writeRows(Arrays.asList( + RowFactory.create("k1", "p1", "v1"), + RowFactory.create("k2", "p1", "v2"), + RowFactory.create("k3", "p1", "v3")), + simpleSchema(), options, basePath(), SaveMode.Overwrite); + + assertMetaColumnPopulation(basePath(), MetaFieldsMode.FILE_NAME_ONLY); + } + + @Test + void clusteringPreservesCommitTimeAndFileNameMode() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + options.put("hoodie.clustering.inline", "true"); + options.put("hoodie.clustering.inline.max.commits", "1"); + options.put("hoodie.clustering.plan.strategy.target.file.max.bytes", "10485760"); + options.put("hoodie.clustering.plan.strategy.small.file.limit", "10485760"); + + writeRows(Arrays.asList( + RowFactory.create("k1", "p1", "v1"), + RowFactory.create("k2", "p1", "v2"), + RowFactory.create("k3", "p1", "v3")), + simpleSchema(), options, basePath(), SaveMode.Overwrite); + + assertMetaColumnPopulation(basePath(), MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME); + } + + @Test + void clusteringPreservesAllMode() { + Map options = baseOptions(); + // ALL is the default. + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + options.put("hoodie.clustering.inline", "true"); + options.put("hoodie.clustering.inline.max.commits", "1"); + options.put("hoodie.clustering.plan.strategy.target.file.max.bytes", "10485760"); + options.put("hoodie.clustering.plan.strategy.small.file.limit", "10485760"); + + writeRows(Arrays.asList( + RowFactory.create("k1", "p1", "v1"), + RowFactory.create("k2", "p1", "v2"), + RowFactory.create("k3", "p1", "v3")), + simpleSchema(), options, basePath(), SaveMode.Overwrite); + + assertMetaColumnPopulation(basePath(), MetaFieldsMode.ALL); + } + + @Test + void clusteringPreservesNoneMode() { + Map options = baseOptions(); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + options.put("hoodie.clustering.inline", "true"); + options.put("hoodie.clustering.inline.max.commits", "1"); + options.put("hoodie.clustering.plan.strategy.target.file.max.bytes", "10485760"); + options.put("hoodie.clustering.plan.strategy.small.file.limit", "10485760"); + + writeRows(Arrays.asList( + RowFactory.create("k1", "p1", "v1"), + RowFactory.create("k2", "p1", "v2"), + RowFactory.create("k3", "p1", "v3")), + simpleSchema(), options, basePath(), SaveMode.Overwrite); + + assertMetaColumnPopulation(basePath(), MetaFieldsMode.NONE); + } + + @Test + void morWithSelectiveModeIsRejected() { + Map options = baseOptions(); + options.put(DataSourceWriteOptions.TABLE_TYPE().key(), "MERGE_ON_READ"); + options.put(HoodieTableConfig.POPULATE_META_FIELDS.key(), "false"); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), MetaFieldsMode.COMMIT_TIME_ONLY.name()); + options.put(DataSourceWriteOptions.OPERATION().key(), DataSourceWriteOptions.BULK_INSERT_OPERATION_OPT_VAL()); + + Throwable thrown = assertThrows(Throwable.class, () -> + writeRows(Collections.singletonList(RowFactory.create("k1", "p1", "v1")), + simpleSchema(), options, basePath(), SaveMode.Overwrite)); + + String rootMessage = rootMessageOf(thrown); + assertTrue(rootMessage.contains("COPY_ON_WRITE") || rootMessage.contains("MoR") || rootMessage.contains("MERGE_ON_READ"), + "Expected MoR-restriction error, got: " + rootMessage); + } + + private static String rootMessageOf(Throwable thrown) { + Throwable root = thrown; + while (root.getCause() != null) { + root = root.getCause(); + } + return root.getMessage() == null ? "" : root.getMessage(); + } +} diff --git a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestHoodieSparkSqlWriter.scala b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestHoodieSparkSqlWriter.scala index 8968caaf348f0..cd1bf55f7230d 100644 --- a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestHoodieSparkSqlWriter.scala +++ b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestHoodieSparkSqlWriter.scala @@ -107,7 +107,7 @@ class TestHoodieSparkSqlWriter extends HoodieSparkWriterTestBase { // fetch all records from parquet files generated from write to hudi val actualDf = sqlContext.read.parquet(fullPartitionPaths(0), fullPartitionPaths(1), fullPartitionPaths(2)) if (!populateMetaFields) { - List(0, 1, 2, 3, 4).foreach(i => assertEquals(0, actualDf.select(HoodieRecord.HOODIE_META_COLUMNS.get(i)).filter(entry => !(entry.mkString(",").equals(""))).count())) + List(0, 1, 2, 3, 4).foreach(i => assertEquals(0, actualDf.select(HoodieRecord.HOODIE_META_COLUMNS.get(i)).filter(entry => !entry.isNullAt(0) && entry.getString(0).nonEmpty).count())) } // remove metadata columns so that expected and actual DFs can be compared as is val trimmedDf = dropMetaFields(actualDf) diff --git a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestHoodieSparkSqlWriterWithTestFormat.scala b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestHoodieSparkSqlWriterWithTestFormat.scala index 2cf0dc20ad08e..4fa26b3993fbb 100644 --- a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestHoodieSparkSqlWriterWithTestFormat.scala +++ b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestHoodieSparkSqlWriterWithTestFormat.scala @@ -106,7 +106,7 @@ class TestHoodieSparkSqlWriterWithTestFormat extends HoodieSparkWriterTestBase { // fetch all records from parquet files generated from write to hudi val actualDf = sqlContext.read.parquet(fullPartitionPaths(0), fullPartitionPaths(1), fullPartitionPaths(2)) if (!populateMetaFields) { - List(0, 1, 2, 3, 4).foreach(i => assertEquals(0, actualDf.select(HoodieRecord.HOODIE_META_COLUMNS.get(i)).filter(entry => !(entry.mkString(",").equals(""))).count())) + List(0, 1, 2, 3, 4).foreach(i => assertEquals(0, actualDf.select(HoodieRecord.HOODIE_META_COLUMNS.get(i)).filter(entry => !entry.isNullAt(0) && entry.getString(0).nonEmpty).count())) } // remove metadata columns so that expected and actual DFs can be compared as is val trimmedDf = dropMetaFields(actualDf) diff --git a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/BootstrapExecutor.java b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/BootstrapExecutor.java index 24ce507b34a7c..5b430b16dc336 100644 --- a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/BootstrapExecutor.java +++ b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/BootstrapExecutor.java @@ -211,6 +211,8 @@ private void initializeTable() throws IOException { .setTableFormat(props.getString(HoodieTableConfig.TABLE_FORMAT.key(), HoodieTableConfig.TABLE_FORMAT.defaultValue())) .setPopulateMetaFields(props.getBoolean( POPULATE_META_FIELDS.key(), POPULATE_META_FIELDS.defaultValue())) + .setMetaFieldsModeFromString(props.getString( + HoodieTableConfig.META_FIELDS_MODE.key(), HoodieTableConfig.META_FIELDS_MODE.defaultValue())) .setArchiveLogFolder(props.getString( TIMELINE_HISTORY_PATH.key(), TIMELINE_HISTORY_PATH.defaultValue())) .setPayloadClassName(cfg.payloadClassName) diff --git a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/StreamSync.java b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/StreamSync.java index 1ffe12f668774..89538e8cfd969 100644 --- a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/StreamSync.java +++ b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/StreamSync.java @@ -478,6 +478,8 @@ HoodieTableMetaClient initializeEmptyTable(HoodieTableMetaClient.TableBuilder ta .setRecordKeyFields(props.getProperty(DataSourceWriteOptions.RECORDKEY_FIELD().key())) .setPopulateMetaFields(props.getBoolean(HoodieTableConfig.POPULATE_META_FIELDS.key(), HoodieTableConfig.POPULATE_META_FIELDS.defaultValue())) + .setMetaFieldsModeFromString(props.getString(HoodieTableConfig.META_FIELDS_MODE.key(), + HoodieTableConfig.META_FIELDS_MODE.defaultValue())) .setKeyGeneratorClassProp(keyGenClassName) .setPartitionValueExtractorClass(partitionValueExtractorClassName) .setOrderingFields(cfg.sourceOrderingFields) diff --git a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieStreamerMetaFieldsMode.java b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieStreamerMetaFieldsMode.java new file mode 100644 index 0000000000000..ee073d2da0095 --- /dev/null +++ b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieStreamerMetaFieldsMode.java @@ -0,0 +1,139 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hudi.utilities.deltastreamer; + +import org.apache.hudi.common.model.HoodieRecord; +import org.apache.hudi.common.model.MetaFieldsMode; +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.testutils.HoodieTestUtils; + +import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Row; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * End-to-end coverage for {@code hoodie.meta.fields.mode} through the HoodieStreamer entrypoint. + * Each parameterized invocation runs a single ingest cycle in the given {@link MetaFieldsMode} and + * verifies both the persisted table property and the actual on-disk parquet column population. + * + *

Rejection paths (unknown token, populate=true+mode, MoR+mode) are exercised in the datasource + * test {@code TestMetaFieldsMode}; this fixture focuses on the streamer control-flow. + */ +public class TestHoodieStreamerMetaFieldsMode extends HoodieDeltaStreamerTestBase { + + @ParameterizedTest + @EnumSource(value = MetaFieldsMode.class, + names = {"ALL", "NONE", "COMMIT_TIME_ONLY", "FILE_NAME_ONLY", "COMMIT_TIME_AND_FILE_NAME"}) + public void testStreamerRespectsMetaFieldsMode(MetaFieldsMode mode) throws Exception { + String tablePath = basePath + "/streamer_meta_fields_mode_" + mode.name(); + HoodieDeltaStreamer.Config cfg = TestHelpers.makeConfig(tablePath, WriteOperationType.INSERT); + // Force CoW; selective modes are CoW-only until MoR log-write is wired. + cfg.tableType = "COPY_ON_WRITE"; + switch (mode) { + case ALL: + // default; nothing to add + break; + case NONE: + cfg.configs.add(HoodieTableConfig.POPULATE_META_FIELDS.key() + "=false"); + break; + default: + cfg.configs.add(HoodieTableConfig.POPULATE_META_FIELDS.key() + "=false"); + cfg.configs.add(HoodieTableConfig.META_FIELDS_MODE.key() + "=" + mode.name()); + break; + } + HoodieDeltaStreamer streamer = new HoodieDeltaStreamer(cfg, jsc); + streamer.getIngestionService().ingestOnce(); + streamer.shutdownGracefully(); + + HoodieTableMetaClient metaClient = HoodieTestUtils.createMetaClient(context, tablePath); + assertEquals(mode, metaClient.getTableConfig().getMetaFieldsMode(), + "streamer must persist mode=" + mode + " on hoodie.properties"); + assertOnDiskMetaColumns(tablePath, mode); + } + + @Test + public void testStreamerRejectsMorWithSelectiveMode() throws Exception { + String tablePath = basePath + "/streamer_mor_selective_rejected"; + HoodieDeltaStreamer.Config cfg = TestHelpers.makeConfig(tablePath, WriteOperationType.BULK_INSERT); + cfg.tableType = "MERGE_ON_READ"; + cfg.configs.add(HoodieTableConfig.POPULATE_META_FIELDS.key() + "=false"); + cfg.configs.add(HoodieTableConfig.META_FIELDS_MODE.key() + "=" + MetaFieldsMode.COMMIT_TIME_ONLY.name()); + + Throwable thrown = assertThrows(Throwable.class, () -> { + HoodieDeltaStreamer streamer = new HoodieDeltaStreamer(cfg, jsc); + streamer.getIngestionService().ingestOnce(); + streamer.shutdownGracefully(); + }); + + String rootMessage = rootMessageOf(thrown); + assertTrue(rootMessage.contains("COPY_ON_WRITE") || rootMessage.contains("MERGE_ON_READ") + || rootMessage.contains("MoR") || rootMessage.contains(HoodieTableConfig.META_FIELDS_MODE.key()), + "Expected MoR-restriction error, got: " + rootMessage); + } + + private void assertOnDiskMetaColumns(String tablePath, MetaFieldsMode expectedMode) { + // Default HoodieTestDataGenerator partitions are YYYY/MM/DD (three levels). + Dataset raw = sparkSession.read().parquet(tablePath + "/*/*/*/*.parquet"); + Row first = raw.select( + HoodieRecord.COMMIT_TIME_METADATA_FIELD, + HoodieRecord.COMMIT_SEQNO_METADATA_FIELD, + HoodieRecord.RECORD_KEY_METADATA_FIELD, + HoodieRecord.PARTITION_PATH_METADATA_FIELD, + HoodieRecord.FILENAME_METADATA_FIELD).first(); + + if (expectedMode.isCommitTimePopulated()) { + assertNotNull(first.get(0), "commit_time must be populated in mode " + expectedMode); + } else { + assertNull(first.get(0), "commit_time must be null in mode " + expectedMode); + } + if (expectedMode.isFileNamePopulated()) { + assertNotNull(first.get(4), "file_name must be populated in mode " + expectedMode); + } else { + assertNull(first.get(4), "file_name must be null in mode " + expectedMode); + } + if (expectedMode == MetaFieldsMode.ALL) { + assertNotNull(first.get(1), "commit_seq_no must be populated in ALL mode"); + assertNotNull(first.get(2), "record_key must be populated in ALL mode"); + assertNotNull(first.get(3), "partition_path must be populated in ALL mode"); + } else { + assertNull(first.get(1), "commit_seq_no must be null outside ALL mode"); + assertNull(first.get(2), "record_key must be null outside ALL mode"); + assertNull(first.get(3), "partition_path must be null outside ALL mode"); + } + } + + private static String rootMessageOf(Throwable thrown) { + Throwable root = thrown; + while (root.getCause() != null) { + root = root.getCause(); + } + return root.getMessage() == null ? "" : root.getMessage(); + } +}