wombatu-kun commented on code in PR #19875:
URL: https://github.com/apache/hudi/pull/19875#discussion_r3975661434
##########
hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieDropPartitionsTool.java:
##########
@@ -357,12 +362,12 @@ private HiveSyncConfig buildHiveSyncProps() {
props.put(DataSourceWriteOptions.HIVE_PASS().key(), cfg.hivePassWord);
props.put(DataSourceWriteOptions.HIVE_URL().key(), cfg.hiveURL);
props.put(DataSourceWriteOptions.HIVE_PARTITION_FIELDS().key(),
cfg.hivePartitionsField);
Review Comment:
With `--hive-partition-field` at its default empty string,
`META_SYNC_PARTITION_FIELDS` is empty and both `HiveSyncTool.syncAllPartitions`
and `syncPartitions` return early, so a `--sync-hive-meta` run without that
flag never drops a partition from the metastore. Requiring it in
`verifyHiveConfigs` alongside the database and table names would surface that
before the drop.
##########
hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieDropPartitionsTool.java:
##########
@@ -301,6 +301,11 @@ public void run() {
log.info(cfg.toString());
Mode mode = Mode.valueOf(cfg.runningMode.toUpperCase());
+ if (cfg.syncToHive) {
Review Comment:
This runs for `dry_run` too, which never reaches `syncToHiveIfNecessary`, so
`--mode dry_run --sync-hive-meta` without a hive database now fails where it
used to print the listing. Was that intended, or should the check sit at the
top of the `DELETE` arm?
##########
hudi-utilities/src/test/java/org/apache/hudi/utilities/TestTableSizeStats.java:
##########
@@ -0,0 +1,296 @@
+/*
+ * 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;
+
+import org.apache.hudi.client.SparkRDDWriteClient;
+import org.apache.hudi.client.WriteClientTestUtils;
+import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.exception.HoodieIOException;
+import org.apache.hudi.testutils.HoodieSparkClientTestBase;
+import org.apache.hudi.utilities.testutils.CapturingLogAppender;
+
+import org.apache.spark.api.java.JavaRDD;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.Set;
+import java.util.UUID;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
+
+import static
org.apache.hudi.common.testutils.HoodieTestDataGenerator.DEFAULT_FIRST_PARTITION_PATH;
+import static
org.apache.hudi.common.testutils.HoodieTestDataGenerator.DEFAULT_SECOND_PARTITION_PATH;
+import static
org.apache.hudi.common.testutils.HoodieTestDataGenerator.DEFAULT_THIRD_PARTITION_PATH;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests {@link TableSizeStats}. The tool reports through its log, so the
assertions read back the lines it
+ * logged for the table and for every partition it decided to include.
+ */
+public class TestTableSizeStats extends HoodieSparkClientTestBase {
+
+ private static final int RECORDS_PER_PARTITION = 4;
+ private static final String PARTITION_STATS_PREFIX = "Partition stats [name:
";
+
+ private static Stream<Arguments> dateIntervalArgs() {
+ return Stream.of(
+ // only the 2016 partition is on or after the start date
+ Arguments.of("2016/1/1", null, 0L,
Collections.singletonList(DEFAULT_FIRST_PARTITION_PATH)),
+ // only the 2015 partitions are before the end date
+ Arguments.of(null, "2016/1/1", 0L,
+ Arrays.asList(DEFAULT_SECOND_PARTITION_PATH,
DEFAULT_THIRD_PARTITION_PATH)),
+ // half open interval [start, end): the start date is included, the
end date is not
+ Arguments.of("2015/3/16", "2015/3/17", 0L,
Collections.singletonList(DEFAULT_SECOND_PARTITION_PATH)),
+ // --num-days walks back from today, so every partition of this table
is out of the window
+ Arguments.of(null, null, 10L, Collections.emptyList()));
Review Comment:
Every partition here is a decade old, so this arm passes for any window
arithmetic - `minusDays`, `minusMonths` or ignoring `--num-days` entirely all
give the empty set. A partition written a day ago, expected to be counted,
would make the arm discriminate.
##########
hudi-utilities/src/test/java/org/apache/hudi/utilities/TestHoodieDropPartitionsTool.java:
##########
@@ -0,0 +1,317 @@
+/*
+ * 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;
+
+import org.apache.hudi.client.SparkRDDWriteClient;
+import org.apache.hudi.client.WriteClientTestUtils;
+import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.model.HoodieBaseFile;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.model.HoodieReplaceCommitMetadata;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.timeline.HoodieInstant;
+import org.apache.hudi.common.table.view.FileSystemViewManager;
+import org.apache.hudi.common.table.view.HoodieTableFileSystemView;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.hive.HoodieHiveSyncException;
+import org.apache.hudi.testutils.HoodieSparkClientTestBase;
+import org.apache.hudi.utilities.testutils.CapturingLogAppender;
+
+import org.apache.spark.api.java.JavaRDD;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.stream.Collectors;
+
+import static
org.apache.hudi.common.testutils.HoodieTestDataGenerator.DEFAULT_FIRST_PARTITION_PATH;
+import static
org.apache.hudi.common.testutils.HoodieTestDataGenerator.DEFAULT_SECOND_PARTITION_PATH;
+import static
org.apache.hudi.common.testutils.HoodieTestDataGenerator.DEFAULT_THIRD_PARTITION_PATH;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests {@link HoodieDropPartitionsTool} against a small three-partition COW
table.
+ */
+public class TestHoodieDropPartitionsTool extends HoodieSparkClientTestBase {
+
+ private static final int RECORDS_PER_PARTITION = 4;
+
+ private HoodieDropPartitionsTool.Config toolConfig(String mode, String
partitions) {
+ HoodieDropPartitionsTool.Config cfg = new
HoodieDropPartitionsTool.Config();
+ cfg.basePath = basePath;
+ cfg.tableName = metaClient.getTableConfig().getTableName();
+ cfg.runningMode = mode;
+ cfg.partitions = partitions;
+ cfg.parallelism = 2;
+ cfg.configs.add(HoodieWriteConfig.TBL_NAME.key() + "=" + cfg.tableName);
+ return cfg;
+ }
+
+ /**
+ * Writes two insert commits: the first spreads records over all three
partitions, the second adds a
+ * second file slice to the first partition.
+ */
+ private void writeThreePartitionTable() {
+ HoodieWriteConfig writeConfig = getConfigBuilder().build();
+ try (SparkRDDWriteClient client = getHoodieWriteClient(writeConfig)) {
+ String firstCommit = WriteClientTestUtils.createNewInstantTime();
+ List<HoodieRecord> firstBatch = new ArrayList<>();
+ for (String partition : Arrays.asList(
+ DEFAULT_FIRST_PARTITION_PATH, DEFAULT_SECOND_PARTITION_PATH,
DEFAULT_THIRD_PARTITION_PATH)) {
+ firstBatch.addAll(dataGen.generateInsertsForPartition(firstCommit,
RECORDS_PER_PARTITION, partition));
+ }
+ writeBatchAndCommit(client, firstCommit, firstBatch);
+
+ String secondCommit = WriteClientTestUtils.createNewInstantTime();
+ writeBatchAndCommit(client, secondCommit,
+ dataGen.generateInsertsForPartition(secondCommit,
RECORDS_PER_PARTITION, DEFAULT_FIRST_PARTITION_PATH));
+ }
+ }
+
+ private void writeBatchAndCommit(SparkRDDWriteClient client, String
instantTime, List<HoodieRecord> records) {
+ WriteClientTestUtils.startCommitWithTime(client, instantTime);
+ JavaRDD<WriteStatus> writeStatuses =
client.insert(jsc.parallelize(records, 1), instantTime);
+ client.commit(instantTime, writeStatuses);
+ }
+
+ private long latestBaseFileCount(String partition) {
+ HoodieTableMetaClient reloaded = HoodieTableMetaClient.reload(metaClient);
+ try (HoodieTableFileSystemView fsView =
FileSystemViewManager.createInMemoryFileSystemView(
+ context, reloaded,
HoodieMetadataConfig.newBuilder().enable(false).build())) {
+ return fsView.getLatestBaseFiles(partition).count();
+ }
+ }
+
+ private List<String> latestFileIds(String partition) {
+ HoodieTableMetaClient reloaded = HoodieTableMetaClient.reload(metaClient);
+ try (HoodieTableFileSystemView fsView =
FileSystemViewManager.createInMemoryFileSystemView(
+ context, reloaded,
HoodieMetadataConfig.newBuilder().enable(false).build())) {
+ return
fsView.getLatestBaseFiles(partition).map(HoodieBaseFile::getFileId).collect(Collectors.toList());
+ }
+ }
+
+ private List<String> completedInstants() {
+ return
HoodieTableMetaClient.reload(metaClient).getActiveTimeline().filterCompletedInstants()
+
.getInstantsAsStream().map(HoodieInstant::requestedTime).collect(Collectors.toList());
+ }
+
+ @Test
+ public void
testDryRunReportsTheFilesItWouldDeleteAndLeavesTheTableUntouched() {
+ writeThreePartitionTable();
+ List<String> instantsBefore = completedInstants();
+ // what the tool prints must be the file ids the two partitions really hold
+ Set<String> expectedReport = new HashSet<>(Arrays.asList(
+ "Partitions : " + DEFAULT_FIRST_PARTITION_PATH + ", corresponding data
file IDs : "
+ + latestFileIds(DEFAULT_FIRST_PARTITION_PATH),
+ "Partitions : " + DEFAULT_SECOND_PARTITION_PATH + ", corresponding
data file IDs : "
+ + latestFileIds(DEFAULT_SECOND_PARTITION_PATH)));
+
+ HoodieDropPartitionsTool.Config cfg = toolConfig("dry_run",
+ DEFAULT_FIRST_PARTITION_PATH + "," + DEFAULT_SECOND_PARTITION_PATH);
+ List<String> messages;
+ try (CapturingLogAppender logs =
CapturingLogAppender.attachTo(HoodieDropPartitionsTool.class)) {
+ new HoodieDropPartitionsTool(jsc, cfg).run();
+ messages = logs.messages();
+ }
+
+ assertTrue(messages.contains("Data files and partitions to delete : "),
messages.toString());
+ assertEquals(expectedReport,
+ messages.stream().filter(m -> m.startsWith("Partitions :
")).collect(Collectors.toSet()));
+ assertTrue(messages.stream().noneMatch(m ->
m.contains(DEFAULT_THIRD_PARTITION_PATH)),
+ "the partition that was not named must not be reported: " + messages);
+
+ assertEquals(instantsBefore, completedInstants(), "dry run must not add
any instant");
+ assertEquals(1, latestBaseFileCount(DEFAULT_FIRST_PARTITION_PATH));
+ assertEquals(1, latestBaseFileCount(DEFAULT_SECOND_PARTITION_PATH));
+ assertEquals(1, latestBaseFileCount(DEFAULT_THIRD_PARTITION_PATH));
+ }
+
+ @Test
+ public void testDeleteMasksOnlyTheRequestedPartitions() throws IOException {
+ writeThreePartitionTable();
+ int instantsBefore = completedInstants().size();
+
+ HoodieDropPartitionsTool.Config cfg = toolConfig("delete",
+ DEFAULT_FIRST_PARTITION_PATH + "," + DEFAULT_SECOND_PARTITION_PATH);
+ new HoodieDropPartitionsTool(jsc, cfg).run();
+
+ HoodieTableMetaClient reloaded = HoodieTableMetaClient.reload(metaClient);
+ assertEquals(instantsBefore + 1, completedInstants().size(), "delete must
add exactly one instant");
+ HoodieInstant replaceInstant =
reloaded.getActiveTimeline().getCompletedReplaceTimeline().lastInstant().get();
+ HoodieReplaceCommitMetadata replaceMetadata =
+ reloaded.getActiveTimeline().readReplaceCommitMetadata(replaceInstant);
+ assertEquals(
+ new HashSet<>(Arrays.asList(DEFAULT_FIRST_PARTITION_PATH,
DEFAULT_SECOND_PARTITION_PATH)),
+ replaceMetadata.getPartitionToReplaceFileIds().keySet());
+ // the file group of the first partition, written by both commits, is
masked
+ assertEquals(1,
replaceMetadata.getPartitionToReplaceFileIds().get(DEFAULT_FIRST_PARTITION_PATH).size());
+
+ assertEquals(0, latestBaseFileCount(DEFAULT_FIRST_PARTITION_PATH));
+ assertEquals(0, latestBaseFileCount(DEFAULT_SECOND_PARTITION_PATH));
+ assertEquals(1, latestBaseFileCount(DEFAULT_THIRD_PARTITION_PATH),
+ "the partition that was not named must survive");
+ }
+
+ /**
+ * The tool takes its write properties either from --props or from repeated
--hoodie-conf, and only defaults
+ * hoodie.meta.fields.mode from the table when the operator did not name it.
Both sources are checked by asking
+ * for a meta-fields mode the table does not have and expecting the write
config gate to reject it.
+ */
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ public void testWritePropertiesComeFromPropsFileAndHoodieConf(boolean
usePropsFile) throws IOException {
+ writeThreePartitionTable();
+
+ HoodieDropPartitionsTool.Config cfg = toolConfig("dry_run",
DEFAULT_THIRD_PARTITION_PATH);
+ String metaFieldsOverride = "hoodie.meta.fields.mode=NONE";
+ if (usePropsFile) {
+ // the file carries the mode, the --hoodie-conf entry already on the
config carries the table name, so
+ // both sources have to be merged for this run to reach the write config
gate
+ Path propsFile = tempDir.resolve("drop-partitions.properties");
+ Files.write(propsFile, Collections.singletonList(metaFieldsOverride),
StandardCharsets.UTF_8);
+ cfg.propsFilePath = propsFile.toAbsolutePath().toString();
+ } else {
+ cfg.configs.add(metaFieldsOverride);
+ }
+
+ HoodieDropPartitionsTool tool = new HoodieDropPartitionsTool(jsc, cfg);
+ Throwable thrown = assertThrows(HoodieException.class, tool::run);
+ assertTrue(stackMessages(thrown).contains("hoodie.meta.fields.mode"),
+ "expected the meta fields mode from the config source to reach the
write config, got: " + thrown);
+ }
+
+ @Test
+ public void testUnsupportedModeFails() {
+ writeThreePartitionTable();
+ HoodieDropPartitionsTool.Config cfg = toolConfig("purge",
DEFAULT_THIRD_PARTITION_PATH);
+ HoodieDropPartitionsTool tool = new HoodieDropPartitionsTool(jsc, cfg);
+
+ HoodieException thrown = assertThrows(HoodieException.class, tool::run);
+ assertTrue(thrown.getMessage().contains("Unable to delete table partitions
in " + basePath));
+ assertTrue(thrown.getCause() instanceof IllegalArgumentException, "got " +
thrown.getCause());
+ assertEquals(0,
HoodieTableMetaClient.reload(metaClient).getActiveTimeline()
+ .getCompletedReplaceTimeline().countInstants());
+ }
+
+ /**
+ * A missing --hive-database is caught before the delete runs, so the
partitions are still there afterwards.
+ */
+ @Test
+ public void testHiveSyncConfigIsVerifiedBeforeTheDrop() {
+ writeThreePartitionTable();
+ HoodieDropPartitionsTool.Config cfg = toolConfig("delete",
DEFAULT_THIRD_PARTITION_PATH);
+ cfg.syncToHive = true;
+ cfg.hiveDataBase = null;
+ HoodieDropPartitionsTool tool = new HoodieDropPartitionsTool(jsc, cfg);
+
+ HoodieException thrown = assertThrows(HoodieException.class, tool::run);
+ assertTrue(thrown.getCause() instanceof IllegalArgumentException, "got " +
thrown.getCause());
+ assertTrue(thrown.getCause().getMessage().contains("--hive-database"));
+ assertEquals(0,
HoodieTableMetaClient.reload(metaClient).getActiveTimeline()
+ .getCompletedReplaceTimeline().countInstants(), "nothing may be
dropped once the hive configs are bad");
+ assertEquals(1, latestBaseFileCount(DEFAULT_THIRD_PARTITION_PATH));
+ }
+
+ /**
+ * With the hive configs in place the sync props are built and the sync is
attempted for real; pointing it at a
+ * port nothing listens on keeps the test free of a metastore. The drop is
committed before that attempt, so a
+ * metastore that is down costs the sync, not the partitions.
+ */
+ @Test
+ public void testHiveSyncFailureLeavesTheDropCommitted() {
+ // the tool feeds the FileSystem's hadoop conf into the HiveConf, which is
the only way in for these
+ jsc.hadoopConfiguration().set("hive.metastore.connect.retries", "1");
+ jsc.hadoopConfiguration().set("hive.metastore.client.connect.retry.delay",
"0s");
+ jsc.hadoopConfiguration().set("hive.metastore.failure.retries", "0");
+ writeThreePartitionTable();
+ HoodieDropPartitionsTool.Config cfg = toolConfig("delete",
DEFAULT_THIRD_PARTITION_PATH);
+ cfg.syncToHive = true;
+ cfg.hiveDataBase = "db";
+ cfg.hiveTableName = "tbl";
+ cfg.hivePartitionsField = "partition_path";
+ cfg.hiveHMSUris = "thrift://localhost:1";
+ HoodieDropPartitionsTool tool = new HoodieDropPartitionsTool(jsc, cfg);
+
+ HoodieException thrown = assertThrows(HoodieException.class, tool::run);
+ assertTrue(thrown.getCause() instanceof HoodieHiveSyncException, "got " +
thrown.getCause());
+ assertTrue(stackMessages(thrown).contains("Failed to create
HiveMetaStoreClient"), stackMessages(thrown));
+ assertTrue(stackMessages(thrown).contains("Could not connect to meta store
using any of the URIs provided"),
+ stackMessages(thrown));
+
+ assertEquals(1,
HoodieTableMetaClient.reload(metaClient).getActiveTimeline()
+ .getCompletedReplaceTimeline().countInstants(), "the drop is committed
before hive sync runs");
+ assertEquals(0, latestBaseFileCount(DEFAULT_THIRD_PARTITION_PATH));
+ }
+
+ @Test
+ public void testConfigEqualsHashCodeAndToString() {
+ HoodieDropPartitionsTool.Config left = new
HoodieDropPartitionsTool.Config();
+ left.basePath = "/tmp/table";
+ left.runningMode = "delete";
+ left.tableName = "t1";
+ left.partitions = "p1,p2";
+ left.configs = new ArrayList<>(Collections.singletonList("k=v"));
+
+ HoodieDropPartitionsTool.Config right = new
HoodieDropPartitionsTool.Config();
+ right.basePath = "/tmp/table";
+ right.runningMode = "delete";
+ right.tableName = "t1";
+ right.partitions = "p1,p2";
+ right.configs = new ArrayList<>(Collections.singletonList("k=v"));
+
+ assertEquals(left, left);
+ assertEquals(left, right);
Review Comment:
`HoodieDropPartitionsTool.Config.equals` and `TableSizeStats.Config.equals`
still compare `basePath` with a bare `.equals`, so a default instance NPEs
there, and both new `testConfigEqualsHashCodeAndToString` set `basePath` on
either side so neither reaches it. Worth the same `Objects.equals` conversion
plus default-instance assertion the validator just got - follow-up, not a
blocker.
--
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]