nsivabalan commented on code in PR #19205: URL: https://github.com/apache/hudi/pull/19205#discussion_r3782123409
########## hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestMetaFieldsModeE2E.java: ########## @@ -0,0 +1,830 @@ +/* + * 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.DataSourceReadOptions; +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.common.table.timeline.HoodieInstant; +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.functions; +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 org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; + +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +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 TestMetaFieldsModeE2E 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<String, String> baseOptions() { + Map<String, String> 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<Row> records, StructType schema, Map<String, String> 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<String, String> 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<Row> 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(); Review Comment: Fixed in b0281670. Both helpers now count non-null values across the whole dataset and require the fixture to write more than one row, so the mutation you named -- populate the opted-in columns for the first record only -- can no longer pass. `assertOnDiskMetaColumns` in `TestHoodieStreamerMetaFieldsMode` had the same defect and got the same treatment. Worth being straight about the verification: I do not have a clean local run of these to point at. Every attempt returned `Tests run: 0`, and I only worked out why afterwards -- `TestMetaFieldsModeE2E` carries no `@Tag`, so the `functional-tests` profile excludes it while `unit-tests` runs it, and I had been reaching for the wrong one. When I finally ran it correctly the suite surfaced a genuine bug in my *other* assertion (see the reply on the commit-time thread), which is a fair illustration of your broader point that these tests were not being exercised. ########## hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestMetaFieldsModeE2E.java: ########## @@ -0,0 +1,830 @@ +/* + * 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.DataSourceReadOptions; +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.common.table.timeline.HoodieInstant; +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.functions; +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 org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; + +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +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 TestMetaFieldsModeE2E 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<String, String> baseOptions() { + Map<String, String> 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<Row> records, StructType schema, Map<String, String> 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<String, String> 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<Row> 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<String, String> 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<String, String> 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<String, String> 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<String, String> 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<String, String> 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 explicitlyContradictingTheModeIsRejectedAtTableCreation() { + // A selective mode implies populate.meta.fields=false. Stating the boolean as true alongside it + // is a contradiction, and the user is told rather than having half their request discarded. + // This is the datasource end of the check in HoodieTableMetaClient.TableBuilder. + Map<String, String> 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()); + + Throwable thrown = assertThrows(Throwable.class, () -> + writeSampleAndGetTableConfig(options, basePath())); + + String rootMessage = rootMessageOf(thrown); + assertTrue(rootMessage.contains(HoodieTableConfig.META_FIELDS_MODE.key()) + && rootMessage.contains(HoodieTableConfig.POPULATE_META_FIELDS.key()), + "the error must name both properties so the user knows which to drop, got: " + rootMessage); + } + + @Test + void selectiveModeWithoutTheLegacyBooleanDerivesItAsFalse() { + // The ordinary case: state only the mode. The boolean is derived, never carried through + // verbatim -- a pre-1.3.0 reader ignores the mode property, so leaving populate=true would make + // it treat a selectively-written table as ALL. + Map<String, String> options = baseOptions(); + 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); + assertFalse(tc.populateMetaFields(), + "legacy populate.meta.fields must be derived from the mode"); + } + + @Test + void noneModePersistsLegacyBooleanAsFalse() { + // The unsafe case this invariant protects: an old incremental reader that saw + // populate.meta.fields=true on a NONE table would run against all-null commit times and + // silently return zero rows. Stating only the mode -- the ordinary case -- must derive false. + Map<String, String> options = baseOptions(); + 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<String, String> options = baseOptions(); + 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<String, String> 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<String, String> 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<String, String> 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<String, String> 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<String, String> 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<String, String> 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. + // + // These target 358fbfdd717a, where HoodieRowCreateHandle's selective path copied the *source + // row's* _hoodie_file_name during clustering, leaving records pointing at a file clustering had + // just replaced. asserting only assertNotNull cannot catch that — the stale value is non-null too + // — so the assertion here compares the column against the file actually holding the row. + // + // Only the selective modes are covered: ALL and NONE route through writeRow / + // writeRowNoMetaFields and never enter the branch the fix touched. + // ------------------------------------------------------------------------- + + private Map<String, String> inlineClusteringOptions(MetaFieldsMode mode) { + Map<String, String> options = baseOptions(); + options.put(HoodieTableConfig.META_FIELDS_MODE.key(), mode.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"); + return options; + } + + /** + * Asserts clustering actually ran and that every surviving row's {@code _hoodie_file_name} names + * the file holding it. + * + * <p>Reads through Hudi rather than globbing the parquet directly: after inline clustering the + * pre-clustering file is still on disk (no cleaning has run), so a raw glob would also inspect + * rows that were replaced and are no longer served. + */ + private void assertClusteredFileNamesPointAtTheirOwnFile(String path, MetaFieldsMode mode) { + HoodieTableMetaClient metaClient = + HoodieTableMetaClient.builder().setBasePath(path).setConf(storageConf()).build(); + assertEquals(mode, metaClient.getTableConfig().getMetaFieldsMode()); + assertEquals(1, metaClient.getActiveTimeline().getCompletedReplaceTimeline().countInstants(), + "clustering must have produced a replacecommit, otherwise this test proves nothing"); + + List<Row> rows = spark().read().format("hudi").load(path) + .withColumn("__containing_file", functions.input_file_name()) + .collectAsList(); + assertFalse(rows.isEmpty(), "expected the clustered table to still serve rows"); + + for (Row row : rows) { + String fileName = row.getAs(HoodieRecord.FILENAME_METADATA_FIELD); + String containingFile = row.getAs("__containing_file").toString(); + if (mode.isFileNamePopulated()) { + assertNotNull(fileName, "file name is opted in, so clustered rows must carry one"); + assertTrue(containingFile.endsWith("/" + fileName), + "_hoodie_file_name must name the file holding the row after clustering, not the " + + "pre-clustering file it was read from; got " + fileName + " inside " + containingFile); + } else { + assertNull(fileName, + "file name is not opted in, so clustering must not populate it; got " + fileName); + } + } + } + + @Test + void clusteringWritesTheNewFileNameUnderFileNameOnly() { + Map<String, String> options = inlineClusteringOptions(MetaFieldsMode.FILE_NAME_ONLY); + writeRows(Arrays.asList( + RowFactory.create("k1", "p1", "v1"), + RowFactory.create("k2", "p1", "v2"), + RowFactory.create("k3", "p1", "v3")), + simpleSchema(), options, basePath(), SaveMode.Overwrite); + + assertClusteredFileNamesPointAtTheirOwnFile(basePath(), MetaFieldsMode.FILE_NAME_ONLY); + } + + @Test + void clusteringWritesTheNewFileNameUnderCommitTimeAndFileName() { + Map<String, String> options = inlineClusteringOptions(MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME); + writeRows(Arrays.asList( + RowFactory.create("k1", "p1", "v1"), + RowFactory.create("k2", "p1", "v2"), + RowFactory.create("k3", "p1", "v3")), + simpleSchema(), options, basePath(), SaveMode.Overwrite); + + assertClusteredFileNamesPointAtTheirOwnFile(basePath(), MetaFieldsMode.COMMIT_TIME_AND_FILE_NAME); + } + + @Test + void clusteringLeavesFileNameNullUnderCommitTimeOnly() { Review Comment: Confirmed and addressed in b0281670: `assertClusteredFileNamesPointAtTheirOwnFile` now captures the replacecommit instant and asserts no clustered row carries it, so clustering stamping the replacecommit instead of preserving the original commit time can no longer pass. Your framing was right that this is the same shape as `358fbfdd717a` -- the file-name assertion was strengthened in response to that fix and the commit-time one was not. One thing outstanding, and it is mine: when I finally got `TestMetaFieldsModeE2E` running under the correct profile, this new assertion **fails** in all three clustering tests -- the commit time comes back null when read through `format("hudi")`, while the raw-parquet assertion in the same tests sees it populated. So either my assertion is reading the wrong thing, or the hudi read path is dropping a column it should serve on a `COMMIT_TIME_ONLY` table. I am digging into which, and would rather report that than leave the assertion looking green because it never ran. Will follow up on this thread. ########## 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()) { Review Comment: Confirmed and fixed in a6d34815. I verified all three parts: `IncrementalRelationV1/V2` are constructed only by `HoodieStreamSourceV1:186` and `HoodieStreamSourceV2:162`, `grep -c readStream TestMetaFieldsModeE2E.java` was 0, and the guard I relaxed sits on exactly that path. Added `streamingReadWorksUnderCommitTimeOnly` (asserts rows come back) and `streamingReadIsRejectedOnMorWithoutMetaFields` (asserts the MoR guard still fires). Your point about where the coverage was missing is well taken -- relaxing a guard on a read path that had already broken twice under selective modes, with nothing exercising the relaxed branch, is not a good place to have no test. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
