This is an automated email from the ASF dual-hosted git repository.
amaliujia pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new 2fdaf29ee9d HDDS-14946. Add CLI command for
RewriteTablePathOzoneAction (#10409)
2fdaf29ee9d is described below
commit 2fdaf29ee9ddf823050d3a5670f479cc5c2037c3
Author: sreejasahithi <[email protected]>
AuthorDate: Wed Jun 10 14:00:43 2026 +0530
HDDS-14946. Add CLI command for RewriteTablePathOzoneAction (#10409)
Co-authored-by: Sreeja Chintalapati
<[email protected]>
---
hadoop-ozone/dist/src/shell/ozone/ozone | 26 ++++-
hadoop-ozone/iceberg/pom.xml | 35 ++++--
.../hadoop/ozone/iceberg/IcebergCommand.java | 42 +++++++
.../ozone/iceberg/RewriteTablePathCommand.java | 127 +++++++++++++++++++++
.../ozone/iceberg/RewriteTablePathOzoneAction.java | 22 ++--
.../iceberg/TestRewriteTablePathOzoneAction.java | 126 ++++++++++++++------
6 files changed, 314 insertions(+), 64 deletions(-)
diff --git a/hadoop-ozone/dist/src/shell/ozone/ozone
b/hadoop-ozone/dist/src/shell/ozone/ozone
index 432ed1d173d..a6c37397654 100755
--- a/hadoop-ozone/dist/src/shell/ozone/ozone
+++ b/hadoop-ozone/dist/src/shell/ozone/ozone
@@ -23,6 +23,16 @@ MYNAME="${BASH_SOURCE-$0}"
bin=$(cd -P -- "$(dirname -- "${MYNAME}")" >/dev/null && pwd -P)
JVM_PID="$$"
+## @description true when the ozone-iceberg artifact is present in this
distribution
+## @audience private
+function ozone_iceberg_available
+{
+ local lib_dir="${HDDS_LIB_JARS_DIR:-${OZONE_HOME}/share/ozone/lib}"
+
+ [[ -f "${OZONE_HOME}/share/ozone/classpath/ozone-iceberg.classpath" ]] \
+ && compgen -G "${lib_dir}/ozone-iceberg-*.jar" > /dev/null
+}
+
## @description build up the ozone command's usage text.
## @audience public
## @stability stable
@@ -65,7 +75,9 @@ function ozone_usage
ozone_add_subcommand "repair" client "Ozone repair tool"
ozone_add_subcommand "ratis" client "Ozone ratis tool"
ozone_add_subcommand "vapor" client "Ozone server simulator"
-
+ if ozone_iceberg_available; then
+ ozone_add_subcommand "iceberg" client "commands for Iceberg tables on
Ozone (see ozone iceberg --help for subcommands)"
+ fi
ozone_generate_usage "${OZONE_SHELL_EXECNAME}" false
}
@@ -252,6 +264,15 @@ function ozonecmd_case
OZONE_VAPOR_OPTS="${OZONE_VAPOR_OPTS} ${RATIS_OPTS}
${OZONE_MODULE_ACCESS_ARGS}"
OZONE_RUN_ARTIFACT_NAME="ozone-vapor"
;;
+ iceberg)
+ if ! ozone_iceberg_available; then
+ ozone_error "ERROR: ozone iceberg is not available in this
distribution (requires JDK 11+ build)."
+ exit 1
+ fi
+ OZONE_CLASSNAME="org.apache.hadoop.ozone.iceberg.IcebergCommand"
+ OZONE_RUN_ARTIFACT_NAME="ozone-iceberg"
+ OZONE_SUBCMD_SUPPORTDAEMONIZATION=false
+ ;;
*)
OZONE_CLASSNAME="${subcmd}"
if ! ozone_validate_classname "${OZONE_CLASSNAME}"; then
@@ -288,7 +309,8 @@ function ozone_suppress_shell_log
{
if [[ "${OZONE_RUN_ARTIFACT_NAME}" =~ ozone-cli-.* ]] \
|| [[ "${OZONE_RUN_ARTIFACT_NAME}" == "ozone-dist" ]] \
- || [[ "${OZONE_RUN_ARTIFACT_NAME}" == "ozone-tools" ]]; then
+ || [[ "${OZONE_RUN_ARTIFACT_NAME}" == "ozone-tools" ]] \
+ || [[ "${OZONE_RUN_ARTIFACT_NAME}" == "ozone-iceberg" ]]; then
if [[ -z "${OZONE_ORIGINAL_LOGLEVEL}" ]] \
&& [[ -z "${OZONE_ORIGINAL_ROOT_LOGGER}" ]]; then
OZONE_LOGLEVEL=OFF
diff --git a/hadoop-ozone/iceberg/pom.xml b/hadoop-ozone/iceberg/pom.xml
index 988f7bc6f66..eea22a04d33 100644
--- a/hadoop-ozone/iceberg/pom.xml
+++ b/hadoop-ozone/iceberg/pom.xml
@@ -32,11 +32,19 @@
</properties>
<dependencies>
+ <dependency>
+ <groupId>info.picocli</groupId>
+ <artifactId>picocli</artifactId>
+ </dependency>
<dependency>
<groupId>org.apache.avro</groupId>
<artifactId>avro</artifactId>
<version>1.12.0</version>
</dependency>
+ <dependency>
+ <groupId>org.apache.hadoop</groupId>
+ <artifactId>hadoop-common</artifactId>
+ </dependency>
<!-- Iceberg dependencies -->
<dependency>
@@ -101,12 +109,20 @@
</exclusion>
</exclusions>
</dependency>
+
+ <dependency>
+ <groupId>org.apache.ozone</groupId>
+ <artifactId>hdds-cli-common</artifactId>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.ozone</groupId>
+ <artifactId>hdds-common</artifactId>
+ </dependency>
<dependency>
<groupId>org.apache.parquet</groupId>
<artifactId>parquet-column</artifactId>
<version>1.16.0</version>
</dependency>
-
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
@@ -175,15 +191,14 @@
</exclusions>
</dependency>
<dependency>
- <groupId>org.apache.hadoop</groupId>
- <artifactId>hadoop-common</artifactId>
- <scope>test</scope>
- <exclusions>
- <exclusion>
- <groupId>org.apache.avro</groupId>
- <artifactId>avro</artifactId>
- </exclusion>
- </exclusions>
+ <groupId>org.apache.ozone</groupId>
+ <artifactId>ozone-filesystem</artifactId>
+ <scope>runtime</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.slf4j</groupId>
+ <artifactId>slf4j-reload4j</artifactId>
+ <scope>runtime</scope>
</dependency>
</dependencies>
diff --git
a/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/IcebergCommand.java
b/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/IcebergCommand.java
new file mode 100644
index 00000000000..af498c07ed9
--- /dev/null
+++
b/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/IcebergCommand.java
@@ -0,0 +1,42 @@
+/*
+ * 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.hadoop.ozone.iceberg;
+
+import org.apache.hadoop.hdds.cli.GenericCli;
+import org.apache.hadoop.hdds.cli.HddsVersionProvider;
+import picocli.CommandLine.Command;
+
+/**
+ * Parent command for Iceberg tables on Ozone.
+ */
+@Command(
+ name = "ozone iceberg",
+ aliases = "iceberg",
+ description = "commands for Iceberg tables on Ozone",
+ subcommands = {
+ RewriteTablePathCommand.class
+ },
+ versionProvider = HddsVersionProvider.class,
+ mixinStandardHelpOptions = true
+)
+public class IcebergCommand extends GenericCli {
+
+ public static void main(String[] args) {
+ new IcebergCommand().run(args);
+ }
+}
diff --git
a/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/RewriteTablePathCommand.java
b/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/RewriteTablePathCommand.java
new file mode 100644
index 00000000000..c8e972750f7
--- /dev/null
+++
b/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/RewriteTablePathCommand.java
@@ -0,0 +1,127 @@
+/*
+ * 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.hadoop.ozone.iceberg;
+
+import java.util.concurrent.Callable;
+import org.apache.hadoop.hdds.cli.AbstractSubcommand;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.actions.RewriteTablePath;
+import org.apache.iceberg.hadoop.HadoopTables;
+import picocli.CommandLine.Command;
+import picocli.CommandLine.Option;
+
+/**
+ * CLI to rewrite Iceberg table paths.
+ */
+@Command(
+ name = "rewrite-path",
+ description = "Rewrite Iceberg table paths for table migration"
+)
+public class RewriteTablePathCommand extends AbstractSubcommand implements
Callable<Void> {
+
+ @Option(
+ names = {"-l", "--table-location"},
+ required = true,
+ description = "The latest metadata.json file path of the table"
+ )
+ private String tableLocation;
+
+ @Option(
+ names = {"-s", "--source-prefix"},
+ required = true,
+ description = "Source path prefix to replace"
+ )
+ private String sourcePrefix;
+
+ @Option(
+ names = {"-t", "--target-prefix"},
+ required = true,
+ description = "Target path prefix"
+ )
+ private String targetPrefix;
+
+ @Option(
+ names = {"--staging"},
+ description = "Staging location where all the rewritten files will be
placed "
+ + "(Default is a new directory under the table's current metadata
directory.)"
+ )
+ private String stagingLocation;
+
+ @Option(
+ names = {"--start-version"},
+ description = "Start version metadata file name (optional, e.g.,
v1.metadata.json)"
+ )
+ private String startVersion;
+
+ @Option(
+ names = {"--end-version"},
+ description = "End version metadata file name (optional, defaults to
current)"
+ )
+ private String endVersion;
+
+ @Option(
+ names = {"--threads"},
+ defaultValue = "10",
+ description = "Number of threads to use (positive integer). "
+ + "If omitted or zero, the default thread count 10 is used."
+ )
+ private int threads;
+
+ @Override
+ public Void call() {
+ out().println("Starting Iceberg table path rewrite");
+ out().println("Table location: " + tableLocation);
+ out().println("Source prefix: " + sourcePrefix);
+ out().println("Target prefix: " + targetPrefix);
+
+ HadoopTables tables = new HadoopTables(getOzoneConf());
+ Table table = tables.load(tableLocation.trim());
+ out().println("Table loaded: " + table.location());
+
+ RewriteTablePathOzoneAction action = new
RewriteTablePathOzoneAction(table, threads);
+ out().println("Threads: " + threads);
+
+ RewriteTablePath rewriteAction =
action.rewriteLocationPrefix(sourcePrefix, targetPrefix);
+
+ if (stagingLocation != null && !stagingLocation.isBlank()) {
+ out().println("Staging location: " + stagingLocation);
+ rewriteAction.stagingLocation(stagingLocation);
+ }
+
+ if (startVersion != null && !startVersion.isBlank()) {
+ out().println("Start version: " + startVersion);
+ rewriteAction.startVersion(startVersion);
+ }
+
+ if (endVersion != null && !endVersion.isBlank()) {
+ out().println("End version: " + endVersion);
+ rewriteAction.endVersion(endVersion);
+ }
+
+ RewriteTablePath.Result result = rewriteAction.execute();
+
+ out().println();
+ out().println("Rewrite completed successfully");
+ out().println(" Latest version: " + result.latestVersion());
+ out().println(" Staging location: " + result.stagingLocation());
+ out().println();
+ out().println("Next step: Copy files from source to target using the file
list");
+ out().println(" File list location: " + result.fileListLocation());
+ return null;
+ }
+}
diff --git
a/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/RewriteTablePathOzoneAction.java
b/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/RewriteTablePathOzoneAction.java
index 4dc50434a3c..4a025b6e935 100644
---
a/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/RewriteTablePathOzoneAction.java
+++
b/hadoop-ozone/iceberg/src/main/java/org/apache/hadoop/ozone/iceberg/RewriteTablePathOzoneAction.java
@@ -96,22 +96,16 @@ public class RewriteTablePathOzoneAction implements
RewriteTablePath {
private String startVersionName;
private String endVersionName;
private String stagingDir;
- private int parallelism;
+ private int threads;
private ExecutorService executorService;
private static final int MAX_INFLIGHT_MULTIPLIER = 4;
- private static final int DEFAULT_THREAD_COUNT = 10;
private final Table table;
- public RewriteTablePathOzoneAction(Table table) {
+ public RewriteTablePathOzoneAction(Table table, int threads) {
this.table = table;
- this.parallelism = DEFAULT_THREAD_COUNT;
- }
-
- public RewriteTablePathOzoneAction(Table table, int parallelism) {
- this.table = table;
- this.parallelism = parallelism;
+ this.threads = threads;
}
@Override
@@ -147,7 +141,7 @@ public RewriteTablePath stagingLocation(String
stagingLocation) {
@Override
public Result execute() {
validateInputs();
- executorService = Executors.newFixedThreadPool(parallelism);
+ executorService = Executors.newFixedThreadPool(threads);
try {
return doExecute();
} finally {
@@ -326,7 +320,7 @@ private Set<Pair<String, String>>
rewriteVersionFile(TableMetadata metadata, Str
private Set<String> manifestsToRewrite(Set<Snapshot> validSnapshots,
Set<Long> deltaSnapshotIds) {
Set<String> manifestPaths = ConcurrentHashMap.newKeySet();
- int maxInFlight = parallelism * MAX_INFLIGHT_MULTIPLIER;
+ int maxInFlight = threads * MAX_INFLIGHT_MULTIPLIER;
Semaphore semaphore = new Semaphore(maxInFlight);
ExecutorCompletionService<Void> completionService = new
ExecutorCompletionService<>(executorService);
@@ -441,7 +435,7 @@ private RewriteResult<ManifestFile>
rewriteManifestLists(Set<Snapshot> validSnap
return new RewriteResult<>();
}
- int maxInFlight = parallelism * MAX_INFLIGHT_MULTIPLIER;
+ int maxInFlight = threads * MAX_INFLIGHT_MULTIPLIER;
Semaphore semaphore = new Semaphore(maxInFlight);
ExecutorCompletionService<RewriteResult<ManifestFile>> completionService =
new ExecutorCompletionService<>(executorService);
@@ -535,7 +529,7 @@ private RewriteContentFileResult rewriteManifests(
return new RewriteContentFileResult();
}
- int maxInFlight = parallelism * MAX_INFLIGHT_MULTIPLIER;
+ int maxInFlight = threads * MAX_INFLIGHT_MULTIPLIER;
Semaphore semaphore = new Semaphore(maxInFlight);
ExecutorCompletionService<RewriteContentFileResult> completionService =
new ExecutorCompletionService<>(executorService);
@@ -735,7 +729,7 @@ private void rewritePositionDeletes(Set<DeleteFile>
toRewrite) {
}
RewriteTablePathUtil.PositionDeleteReaderWriter posDeleteReaderWriter =
new OzonePositionDeleteReaderWriter();
- int maxInFlight = parallelism * MAX_INFLIGHT_MULTIPLIER;
+ int maxInFlight = threads * MAX_INFLIGHT_MULTIPLIER;
Semaphore semaphore = new Semaphore(maxInFlight);
ExecutorCompletionService<Void> completionService = new
ExecutorCompletionService<>(executorService);
int submittedTasks = 0;
diff --git
a/hadoop-ozone/iceberg/src/test/java/org/apache/hadoop/ozone/iceberg/TestRewriteTablePathOzoneAction.java
b/hadoop-ozone/iceberg/src/test/java/org/apache/hadoop/ozone/iceberg/TestRewriteTablePathOzoneAction.java
index 5b11edc1a11..514ec338ec8 100644
---
a/hadoop-ozone/iceberg/src/test/java/org/apache/hadoop/ozone/iceberg/TestRewriteTablePathOzoneAction.java
+++
b/hadoop-ozone/iceberg/src/test/java/org/apache/hadoop/ozone/iceberg/TestRewriteTablePathOzoneAction.java
@@ -23,11 +23,14 @@
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.io.BufferedReader;
+import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.InputStreamReader;
+import java.io.PrintStream;
import java.nio.charset.StandardCharsets;
import java.nio.file.Path;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@@ -35,7 +38,7 @@
import java.util.Set;
import java.util.UUID;
import java.util.stream.Collectors;
-import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.DataFiles;
import org.apache.iceberg.DeleteFile;
@@ -78,6 +81,7 @@
import org.apache.iceberg.parquet.Parquet;
import org.apache.iceberg.types.Types;
import org.apache.iceberg.util.Pair;
+import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -102,6 +106,11 @@ class TestRewriteTablePathOzoneAction {
private String targetPrefix = null;
private Table table = null;
+ private ByteArrayOutputStream outContent;
+ private ByteArrayOutputStream errContent;
+ private PrintStream originalOut;
+ private PrintStream originalErr;
+
@TempDir
private Path tableDir;
@TempDir
@@ -115,14 +124,24 @@ public void setupTableLocation() throws IOException {
this.table = createTable(tableLocation + "/");
this.sourcePrefix = tableLocation;
this.targetPrefix = targetDir.toUri().toString().replaceFirst("^file:///",
"file:/") + TABLE_NAME;
+
+ outContent = new ByteArrayOutputStream();
+ errContent = new ByteArrayOutputStream();
+ originalOut = System.out;
+ originalErr = System.err;
+ System.setOut(new PrintStream(outContent, true, StandardCharsets.UTF_8));
+ System.setErr(new PrintStream(errContent, true, StandardCharsets.UTF_8));
+ }
+
+ @AfterEach
+ public void restoreStreams() {
+ System.setOut(originalOut);
+ System.setErr(originalErr);
}
@Test
void fullTablePathRewrite() throws Exception {
- RewriteTablePath.Result result = new RewriteTablePathOzoneAction(table, 2)
- .rewriteLocationPrefix(sourcePrefix, targetPrefix)
- .stagingLocation(stagingDir.toString() + "/")
- .execute();
+ String fileListLocation = executeRewriteCommand("--threads", "2");
List<String> metadataPaths = metadataLogEntryPaths(table);
Set<String> expectedTargets = new HashSet<>();
@@ -130,7 +149,7 @@ void fullTablePathRewrite() throws Exception {
expectedTargets.add(RewriteTablePathUtil.newPath(path, sourcePrefix,
targetPrefix));
}
- Set<Pair<String, String>> csvPairs = readCsvPairs(table,
result.fileListLocation());
+ Set<Pair<String, String>> csvPairs = readCsvPairs(table, fileListLocation);
Set<String> actualTargets = csvPairs.stream().map(Pair::second)
.filter(p -> p.endsWith(".metadata.json"))
.collect(Collectors.toSet());
@@ -146,11 +165,7 @@ void tablePathRewriteForStartAndNoEndVersionProvided()
throws Exception {
List<String> metadataPaths = metadataLogEntryPaths(table);
String startName = RewriteTablePathUtil.fileName(metadataPaths.get(2));
- RewriteTablePath.Result result = new RewriteTablePathOzoneAction(table)
- .rewriteLocationPrefix(sourcePrefix, targetPrefix)
- .stagingLocation(stagingDir.toString() + "/")
- .startVersion(startName)
- .execute();
+ String fileListLocation = executeRewriteCommand("--start-version",
startName);
List<String> expectedPaths = new ArrayList<>();
for (int i = metadataPaths.size() - 1; i >= 3; i--) {
@@ -162,7 +177,7 @@ void tablePathRewriteForStartAndNoEndVersionProvided()
throws Exception {
expectedTargets.add(RewriteTablePathUtil.newPath(versionPath,
sourcePrefix, targetPrefix));
}
- Set<Pair<String, String>> csvPairs = readCsvPairs(table,
result.fileListLocation());
+ Set<Pair<String, String>> csvPairs = readCsvPairs(table, fileListLocation);
Set<String> actualTargets = csvPairs.stream().map(Pair::second)
.filter(p -> p.endsWith(".metadata.json"))
.collect(Collectors.toSet());
@@ -178,11 +193,7 @@ void tablePathRewriteForOnlyEndVersionProvided() throws
Exception {
List<String> metadataPaths = metadataLogEntryPaths(table);
String endName = RewriteTablePathUtil.fileName(metadataPaths.get(2));
- RewriteTablePath.Result result = new RewriteTablePathOzoneAction(table)
- .rewriteLocationPrefix(sourcePrefix, targetPrefix)
- .stagingLocation(stagingDir.toString() + "/")
- .endVersion(endName)
- .execute();
+ String fileListLocation = executeRewriteCommand("--end-version", endName);
List<String> expectedPaths = new ArrayList<>();
for (int i = 2; i >= 0; i--) {
@@ -194,7 +205,7 @@ void tablePathRewriteForOnlyEndVersionProvided() throws
Exception {
expectedTargets.add(RewriteTablePathUtil.newPath(versionPath,
sourcePrefix, targetPrefix));
}
- Set<Pair<String, String>> csvPairs = readCsvPairs(table,
result.fileListLocation());
+ Set<Pair<String, String>> csvPairs = readCsvPairs(table, fileListLocation);
Set<String> actualTargets = csvPairs.stream().map(Pair::second)
.filter(p -> p.endsWith(".metadata.json"))
.collect(Collectors.toSet());
@@ -211,12 +222,9 @@ void tablePathRewriteForStartAndEndVersionProvided()
throws Exception {
String startName = RewriteTablePathUtil.fileName(metadataPaths.get(1));
String endName = RewriteTablePathUtil.fileName(metadataPaths.get(3));
- RewriteTablePath.Result result = new RewriteTablePathOzoneAction(table)
- .rewriteLocationPrefix(sourcePrefix, targetPrefix)
- .stagingLocation(stagingDir.toString() + "/")
- .startVersion(startName)
- .endVersion(endName)
- .execute();
+ String fileListLocation = executeRewriteCommand(
+ "--start-version", startName,
+ "--end-version", endName);
List<String> expectedPaths = new ArrayList<>();
for (int i = 3; i >= 2; i--) {
@@ -228,7 +236,7 @@ void tablePathRewriteForStartAndEndVersionProvided() throws
Exception {
expectedTargets.add(RewriteTablePathUtil.newPath(versionPath,
sourcePrefix, targetPrefix));
}
- Set<Pair<String, String>> csvPairs = readCsvPairs(table,
result.fileListLocation());
+ Set<Pair<String, String>> csvPairs = readCsvPairs(table, fileListLocation);
Set<String> actualTargets = csvPairs.stream().map(Pair::second)
.filter(p -> p.endsWith(".metadata.json"))
.collect(Collectors.toSet());
@@ -242,7 +250,7 @@ void tablePathRewriteForStartAndEndVersionProvided() throws
Exception {
@Test
void executeRejectsMissingLocationPrefix() {
NullPointerException exception = assertThrows(NullPointerException.class,
- () -> new RewriteTablePathOzoneAction(table)
+ () -> new RewriteTablePathOzoneAction(table, 2)
.stagingLocation(stagingDir.toString() + "/")
.execute());
@@ -252,7 +260,7 @@ void executeRejectsMissingLocationPrefix() {
@Test
void executeRejectsMissingTargetPrefix() {
NullPointerException exception = assertThrows(NullPointerException.class,
- () -> new RewriteTablePathOzoneAction(table)
+ () -> new RewriteTablePathOzoneAction(table, 2)
.rewriteLocationPrefix(sourcePrefix, null));
assertEquals("Target prefix is null", exception.getMessage());
@@ -261,7 +269,7 @@ void executeRejectsMissingTargetPrefix() {
@Test
void rewriteLocationPrefixRejectsSameSourceAndTarget() {
IllegalArgumentException exception =
assertThrows(IllegalArgumentException.class,
- () -> new RewriteTablePathOzoneAction(table)
+ () -> new RewriteTablePathOzoneAction(table, 2)
.rewriteLocationPrefix(sourcePrefix, sourcePrefix)
.execute());
@@ -272,7 +280,7 @@ void rewriteLocationPrefixRejectsSameSourceAndTarget() {
@Test
void startVersionRejectsUnknownVersion() {
IllegalArgumentException exception =
assertThrows(IllegalArgumentException.class,
- () -> new RewriteTablePathOzoneAction(table)
+ () -> new RewriteTablePathOzoneAction(table, 2)
.rewriteLocationPrefix(sourcePrefix, targetPrefix)
.startVersion("missing.metadata.json")
.execute());
@@ -288,7 +296,7 @@ void startVersionRejectsDeletedVersionFile() {
table.io().deleteFile(metadataPaths.get(0));
IllegalArgumentException exception =
assertThrows(IllegalArgumentException.class,
- () -> new RewriteTablePathOzoneAction(table)
+ () -> new RewriteTablePathOzoneAction(table, 2)
.rewriteLocationPrefix(sourcePrefix, targetPrefix)
.startVersion(existingName)
.execute());
@@ -299,7 +307,7 @@ void startVersionRejectsDeletedVersionFile() {
@Test
void endVersionRejectsUnknownVersion() {
IllegalArgumentException exception =
assertThrows(IllegalArgumentException.class,
- () -> new RewriteTablePathOzoneAction(table)
+ () -> new RewriteTablePathOzoneAction(table, 2)
.rewriteLocationPrefix(sourcePrefix, targetPrefix)
.endVersion("missing.metadata.json")
.execute());
@@ -315,7 +323,7 @@ void endVersionRejectsDeletedVersionFile() {
table.io().deleteFile(metadataPaths.get(0));
IllegalArgumentException exception =
assertThrows(IllegalArgumentException.class,
- () -> new RewriteTablePathOzoneAction(table)
+ () -> new RewriteTablePathOzoneAction(table, 2)
.rewriteLocationPrefix(sourcePrefix, targetPrefix)
.endVersion(existingName)
.execute());
@@ -326,7 +334,7 @@ void endVersionRejectsDeletedVersionFile() {
@Test
void usesCurrentMetadataIfEndVersionNotProvided() {
String currentMetadata = ((HasTableOperations)
table).operations().current().metadataFileLocation();
- RewriteTablePathOzoneAction action = new
RewriteTablePathOzoneAction(table);
+ RewriteTablePathOzoneAction action = new
RewriteTablePathOzoneAction(table, 2);
action.rewriteLocationPrefix(sourcePrefix,
targetPrefix).stagingLocation(stagingDir + "/");
RewriteTablePath.Result result = action.execute();
assertThat(result.latestVersion()).isEqualTo(RewriteTablePathUtil.fileName(currentMetadata));
@@ -335,7 +343,7 @@ void usesCurrentMetadataIfEndVersionNotProvided() {
@Test
void defaultStagingDirIsUnderTableMetadataLocation() {
String metadataLocation =
RewriteTablePathOzoneUtils.getMetadataLocation(table);
- RewriteTablePath.Result result = new RewriteTablePathOzoneAction(table)
+ RewriteTablePath.Result result = new RewriteTablePathOzoneAction(table, 2)
.rewriteLocationPrefix(sourcePrefix, targetPrefix)
.execute();
@@ -409,7 +417,7 @@ void rejectsTablesWithPartitionStatistics() {
TableOperations ops = ((HasTableOperations) table).operations();
ops.commit(baseMetadata, metadataWithStats);
- RewriteTablePath action = new RewriteTablePathOzoneAction(table)
+ RewriteTablePath action = new RewriteTablePathOzoneAction(table, 2)
.rewriteLocationPrefix(sourcePrefix, targetPrefix)
.stagingLocation(stagingDir + "/");
@@ -488,7 +496,7 @@ void manifestsToRewriteRejectsMissingManifestList() {
String manifestListLocation = snapshot.manifestListLocation();
table.io().deleteFile(manifestListLocation);
- RewriteTablePath action = new RewriteTablePathOzoneAction(table)
+ RewriteTablePath action = new RewriteTablePathOzoneAction(table, 2)
.rewriteLocationPrefix(sourcePrefix, targetPrefix)
.stagingLocation(stagingDir + "/");
@@ -497,7 +505,49 @@ void manifestsToRewriteRejectsMissingManifestList() {
assertThat(exception.getCause()).hasMessageContaining("Failed to read
manifests for snapshot " +
snapshot.snapshotId());
}
-
+
+ private String executeRewriteCommand(String... optionalArgs) {
+ List<String> args = new ArrayList<>();
+ args.add("rewrite-path");
+ args.add("-l");
+ args.add(table.location());
+ args.add("-s");
+ args.add(sourcePrefix);
+ args.add("-t");
+ args.add(targetPrefix);
+ args.add("--staging");
+ args.add(stagingDir + "/");
+ args.addAll(Arrays.asList(optionalArgs));
+
+ int exitCode = new IcebergCommand().getCmd().execute(args.toArray(new
String[0]));
+ assertEquals(0, exitCode,
+ "Command failed.\nstdout:\n" + stdout() + "\nstderr:\n" + stderr());
+ assertThat(stdout())
+ .contains("Starting Iceberg table path rewrite")
+ .contains("Table loaded: " + table.location())
+ .contains("Staging location: " + stagingDir + "/")
+ .contains("File list location:");
+ return parseFileListLocation(stdout());
+ }
+
+ private String stdout() {
+ return outContent.toString(StandardCharsets.UTF_8);
+ }
+
+ private String stderr() {
+ return errContent.toString(StandardCharsets.UTF_8);
+ }
+
+ private static String parseFileListLocation(String output) {
+ for (String line : output.split("\n")) {
+ if (line.contains("File list location:")) {
+ return line.substring(line.indexOf("File list location:") + "File list
location:".length())
+ .trim();
+ }
+ }
+ throw new IllegalStateException("File list location not found in command
output: " + output);
+ }
+
/**
* For every staged file in the CSV copy plan, asserts that internal paths
are rewritten
* to the target prefix:
@@ -724,7 +774,7 @@ private static Set<Pair<String, String>> readCsvPairs(Table
tbl, String fileList
}
private Table createTable(String location) throws IOException {
- HadoopTables tables = new HadoopTables(new Configuration());
+ HadoopTables tables = new HadoopTables(new OzoneConfiguration());
Table tbl = tables.create(SCHEMA, PartitionSpec.unpartitioned(), new
HashMap<>(), location);
for (int i = 0; i < COMMITS; i++) {
String dataPath = location + "data/batch-" + i + ".parquet";
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]