This is an automated email from the ASF dual-hosted git repository.
szetszwo 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 758883e7bc8 HDDS-16219. Add ExportJob and expand ExportFileManager
(#11052)
758883e7bc8 is described below
commit 758883e7bc8321f96fbd32fc0939b1c96974552d
Author: Sarveksha Yeshavantha Raju
<[email protected]>
AuthorDate: Fri Aug 21 08:27:58 2026 +0530
HDDS-16219. Add ExportJob and expand ExportFileManager (#11052)
---
.../scm/container/export/ExportFileManager.java | 120 ++++++++++++---------
.../hdds/scm/container/export/ExportJob.java | 41 ++++++-
.../hdds/scm/container/export/ExportScope.java | 8 +-
.../container/export/TestExportFileManager.java | 103 +++++++++++++++---
.../hdds/scm/container/export/TestExportJob.java | 67 ++++++++++++
5 files changed, 268 insertions(+), 71 deletions(-)
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportFileManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportFileManager.java
index e0bca340a24..5ebbf012970 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportFileManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportFileManager.java
@@ -17,11 +17,13 @@
package org.apache.hadoop.hdds.scm.container.export;
+import java.io.BufferedWriter;
import java.io.File;
import java.io.IOException;
import java.io.RandomAccessFile;
import java.nio.channels.FileLock;
import java.nio.channels.OverlappingFileLockException;
+import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
@@ -31,7 +33,11 @@
import java.util.Comparator;
import java.util.List;
import java.util.Objects;
+import java.util.zip.GZIPOutputStream;
+import org.apache.commons.compress.archivers.ArchiveOutputStream;
+import org.apache.commons.compress.archivers.tar.TarArchiveEntry;
import org.apache.commons.io.FileUtils;
+import org.apache.hadoop.hdds.utils.Archiver;
import org.apache.hadoop.ozone.util.UUIDUtil;
import org.apache.ratis.util.AtomicFileOutputStream;
import org.slf4j.Logger;
@@ -44,33 +50,28 @@
* uses the layout below. The manager gzip-compresses the archive ({@code
.tar.gz}) so operators
* can stream entries with {@code zcat}.
*
- * <p>While a job runs, shard text files are written under {@code
export_{jobId}/}. The archive is
- * created only after all shards are written. The export manager writes
- * {@code container-ids_{scope}_{timestamp}_job{jobId}.tar.gz.tmp} and
atomically renames it to
+ * <p>While a job runs, part text files are written under {@code
export_{jobId}/}. The archive is
+ * created only after all parts are written. The export manager writes
+ * {@code container-ids_{scope}_{jobStartTime}_job{jobId}.tar.gz.tmp} and
atomically renames it to
* {@code .tar.gz} on close ({@link AtomicFileOutputStream}), so a partial
{@code .tar.gz} is
- * never visible. {@link #lock()} uses {@code in_use.lock} to exclude
concurrent writers.
+ * never visible. {@link #start()} acquires {@code in_use.lock} to exclude
concurrent writers.
*
* <pre>
* {exportDirectory}/
* ├── in_use.lock
- * ├── container-ids_{scope}_{timestamp}_job{jobId}.tar.gz
- * ├── container-ids_{scope}_{timestamp}_job{jobId}.tar.gz.tmp
+ * ├── container-ids_{scope}_{jobStartTime}_job{jobId}.tar.gz
+ * ├── container-ids_{scope}_{jobStartTime}_job{jobId}.tar.gz.tmp
* └── export_{jobId}/
- * ├── container-ids_{scope}_{metadataTimestamp}_part001.txt
+ * ├── container-ids_{scope}_{jobStartTime}_part001.txt
* └── ...
* </pre>
*
* <p><b>Incomplete work</b> ({@code export_{jobId}/} and {@code .tar.gz.tmp})
is removed by
- * {@link #cleanupFailedJob(Path, File)} on failure or cancel, and by {@link
#start()} for every
+ * {@link #cleanupFailedJob(ExportJob.Id, String)} on failure or cancel, and
by {@link #start()} for every
* leftover directory and temp file after SCM restart. Completed {@code
.tar.gz} files are kept.
*
- * <p><b>Completed {@code .tar.gz}</b> remains on disk until the export
manager evicts it
- * ({@code maxTerminalJobs} in {@code ContainerExportManager}) via {@link
#deleteExportTar(String)}.
- *
- * <p><b>SCM restart:</b> in-memory job status is lost. {@link #start()}
clears incomplete work;
- * {@link #listCompletedArchivePaths()} returns existing {@code tarPath}
values (oldest first);
- * {@link #jobIdFromArchiveFileName(String)} parses {@code jobId} for
terminal-job rebuild in
- * {@code ContainerExportManager}.
+ * <p><b>SCM restart:</b> {@link #start()} acquires the export directory lock
and clears incomplete work.
+ * {@link #listCompletedArchivePaths()} returns existing archive paths (oldest
first).
*/
final class ExportFileManager {
@@ -81,7 +82,7 @@ final class ExportFileManager {
static final String EXPORT_ARCHIVE_SUFFIX = ".tar.gz";
static final String EXPORT_ARCHIVE_TMP_SUFFIX = EXPORT_ARCHIVE_SUFFIX +
AtomicFileOutputStream.TMP_EXTENSION;
static final String EXPORT_LOCK_NAME = "in_use.lock";
- private static final int ARCHIVE_TIMESTAMP_LENGTH = 16;
+ private static final int EXPORT_JOB_START_TIME_LENGTH = 19;
private final String exportDirectory;
private FileLock exportDirectoryLock;
@@ -90,16 +91,13 @@ final class ExportFileManager {
this.exportDirectory = Objects.requireNonNull(exportDirectory,
"exportDirectory == null");
}
- String getExportDirectory() {
- return exportDirectory;
- }
-
void start() throws IOException {
Files.createDirectories(Paths.get(exportDirectory));
+ lock();
removeIncompleteWorkOnStartup();
}
- void lock() throws IOException {
+ private void lock() throws IOException {
if (exportDirectoryLock != null) {
return;
}
@@ -119,26 +117,49 @@ void lock() throws IOException {
}
}
- void unlock() throws IOException {
- if (exportDirectoryLock == null) {
- return;
- }
- exportDirectoryLock.release();
- exportDirectoryLock.channel().close();
- exportDirectoryLock = null;
+ File resolveArchiveFile(ExportScope scope, String jobStartTime, ExportJob.Id
jobId) {
+ return new File(exportDirectory, String.format("container-ids_%s_%s%s%s%s",
+ scope.getValue(), jobStartTime, EXPORT_ARCHIVE_JOB_INFIX,
jobId.getValue(), EXPORT_ARCHIVE_SUFFIX));
}
- File resolveArchiveFile(ExportScope scope, String archiveTimestamp,
ExportJob.Id jobId) {
- return new File(exportDirectory, String.format("container-ids_%s_%s%s%s%s",
- scope.getValue(), archiveTimestamp, EXPORT_ARCHIVE_JOB_INFIX,
jobId.getValue(), EXPORT_ARCHIVE_SUFFIX));
+ File resolveArchiveTempFile(ExportScope scope, String jobStartTime,
ExportJob.Id jobId) {
+ return AtomicFileOutputStream.getTemporaryFile(resolveArchiveFile(scope,
jobStartTime, jobId));
+ }
+
+ void createJobDirectory(ExportJob.Id jobId) throws IOException {
+ Files.createDirectories(jobDirectory(jobId));
}
- File resolveArchiveTempFile(ExportScope scope, String archiveTimestamp,
ExportJob.Id jobId) {
- return AtomicFileOutputStream.getTemporaryFile(resolveArchiveFile(scope,
archiveTimestamp, jobId));
+ BufferedWriter newPartWriter(ExportJob.Id jobId, String partFileName) throws
IOException {
+ return Files.newBufferedWriter(jobDirectory(jobId).resolve(partFileName),
StandardCharsets.UTF_8);
+ }
+
+ void writeArchive(ExportJob.Id jobId, String archivePath) throws IOException
{
+ Path jobDir = jobDirectory(jobId);
+ File[] parts = jobDir.toFile().listFiles((dir, name) ->
name.endsWith(".txt"));
+ if (parts == null || parts.length == 0) {
+ throw new IOException("No part files found for export job " + jobId);
+ }
+ Arrays.sort(parts, Comparator.comparing(File::getName));
+ File archiveFile = new File(archivePath);
+ try (AtomicFileOutputStream atomicOut = new
AtomicFileOutputStream(archiveFile);
+ GZIPOutputStream gzipOut = new GZIPOutputStream(atomicOut);
+ ArchiveOutputStream<TarArchiveEntry> tarOut = Archiver.tar(gzipOut)) {
+ for (File part : parts) {
+ Archiver.includeFile(part, part.getName(), tarOut);
+ }
+ }
+ }
+
+ void deleteJobDirectory(ExportJob.Id jobId) throws IOException {
+ Path jobDir = jobDirectory(jobId);
+ if (Files.exists(jobDir)) {
+ FileUtils.deleteDirectory(jobDir.toFile());
+ }
}
/**
- * Returns completed archive paths ({@code tarPath} in {@code
ExportJob.Status}), oldest first.
+ * Returns completed archive paths, oldest first.
*/
List<String> listCompletedArchivePaths() {
File exportDir = new File(exportDirectory);
@@ -148,7 +169,7 @@ List<String> listCompletedArchivePaths() {
return Collections.emptyList();
}
Arrays.sort(matches, Comparator.comparing(
- file -> archiveTimestampFromArchiveFileName(file.getName())));
+ file -> jobStartTimeFromArchiveFileName(file.getName())));
List<String> archivePaths = new ArrayList<>(matches.length);
for (File archive : matches) {
archivePaths.add(archive.getAbsolutePath());
@@ -156,14 +177,14 @@ List<String> listCompletedArchivePaths() {
return archivePaths;
}
- static String archiveTimestampFromArchiveFileName(String fileName) {
+ static String jobStartTimeFromArchiveFileName(String fileName) {
int jobIndex = fileName.lastIndexOf(EXPORT_ARCHIVE_JOB_INFIX);
- if (jobIndex < ARCHIVE_TIMESTAMP_LENGTH + 1
+ if (jobIndex < EXPORT_JOB_START_TIME_LENGTH + 1
|| !fileName.endsWith(EXPORT_ARCHIVE_SUFFIX)
|| fileName.endsWith(EXPORT_ARCHIVE_TMP_SUFFIX)) {
return null;
}
- return fileName.substring(jobIndex - ARCHIVE_TIMESTAMP_LENGTH, jobIndex);
+ return fileName.substring(jobIndex - EXPORT_JOB_START_TIME_LENGTH,
jobIndex);
}
static ExportJob.Id jobIdFromArchiveFileName(String fileName) {
@@ -179,24 +200,17 @@ static ExportJob.Id jobIdFromArchiveFileName(String
fileName) {
return UUIDUtil.isValidUuidString(jobId) ? ExportJob.Id.of(jobId) : null;
}
- void deleteExportTar(String tarPath) {
- if (tarPath == null) {
- return;
- }
- File archive = new File(tarPath);
- if (archive.isFile() && FileUtils.deleteQuietly(archive)) {
- LOG.debug("Removed container export archive: {}", archive.getName());
+ void cleanupFailedJob(ExportJob.Id jobId, String archivePath) {
+ FileUtils.deleteQuietly(jobDirectory(jobId).toFile());
+ if (archivePath != null) {
+ File archive = new File(archivePath);
+ FileUtils.deleteQuietly(archive);
+
FileUtils.deleteQuietly(AtomicFileOutputStream.getTemporaryFile(archive));
}
- FileUtils.deleteQuietly(AtomicFileOutputStream.getTemporaryFile(archive));
}
- void cleanupFailedJob(Path jobDir, File archiveFile) {
- if (jobDir != null) {
- FileUtils.deleteQuietly(jobDir.toFile());
- }
- if (archiveFile != null) {
-
FileUtils.deleteQuietly(AtomicFileOutputStream.getTemporaryFile(archiveFile));
- }
+ private Path jobDirectory(ExportJob.Id jobId) {
+ return Paths.get(exportDirectory, exportJobDirName(jobId));
}
private void removeIncompleteWorkOnStartup() {
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportJob.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportJob.java
index f9dee07eaea..ba66a4907f1 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportJob.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportJob.java
@@ -17,14 +17,20 @@
package org.apache.hadoop.hdds.scm.container.export;
+import java.io.BufferedWriter;
+import java.io.IOException;
import java.util.Objects;
import java.util.UUID;
/**
- * Container ID export job identifier.
+ * Metadata for a container ID export job.
*/
public final class ExportJob {
+ private final Id id;
+ private final ExportScope scope;
+ private final String jobStartTime;
+
/**
* Unique job identifier.
*/
@@ -69,6 +75,37 @@ public int hashCode() {
}
}
- private ExportJob() {
+ ExportJob(Id id, ExportScope scope, String jobStartTime) {
+ this.id = id;
+ this.scope = scope;
+ this.jobStartTime = jobStartTime;
+ }
+
+ String partFileName(int partIndex) {
+ return String.format("container-ids_%s_%s_part%03d.txt",
+ scope.getValue(), jobStartTime, partIndex);
+ }
+
+ void writeMetadataHeader(BufferedWriter writer, int partNumber, long
partStartContainerId)
+ throws IOException {
+ writer.write("# jobId=" + id.getValue());
+ writer.newLine();
+ writer.write("# jobStartTime=" + jobStartTime);
+ writer.newLine();
+ if (scope.getHealthState() != null) {
+ writer.write("# healthState=" + scope.getHealthState().name());
+ writer.newLine();
+ }
+ if (scope.getLifeCycleState() != null) {
+ writer.write("# lifecycleState=" + scope.getLifeCycleState().name());
+ writer.newLine();
+ }
+ writer.write("# startContainerId=" + partStartContainerId);
+ writer.newLine();
+ writer.write("# part=" + partNumber);
+ writer.newLine();
+ writer.write("# format=container-id-per-line");
+ writer.newLine();
+ writer.newLine();
}
}
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportScope.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportScope.java
index 921fdf7f588..c6f174a3299 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportScope.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/export/ExportScope.java
@@ -24,7 +24,7 @@
* Container listing filters for an export job.
* An export job filters containers by {@link ContainerHealthState}, {@link
LifeCycleState} or both.
* Example archive name:
- * {@code
container-ids_health-MISSING_lifecycle-OPEN_20260101T120000Z_job{jobId}.tar.gz}
+ * {@code
container-ids_health-MISSING_lifecycle-OPEN_2026-01-01-12-00-00_job{jobId}.tar.gz}
*/
public final class ExportScope {
@@ -40,6 +40,10 @@ private ExportScope(LifeCycleState lifeCycleState,
ContainerHealthState healthSt
}
public static ExportScope of(LifeCycleState lifeCycleState,
ContainerHealthState healthState) {
+ if (lifeCycleState == null && healthState == null) {
+ throw new IllegalArgumentException("At least one of healthState or
lifecycleState filter is required.");
+ }
+
String health = healthState != null ? healthState.name() : ANY;
String lifecycle = lifeCycleState != null ? lifeCycleState.name() : ANY;
String value = "health-" + health + "_lifecycle-" + lifecycle;
@@ -55,7 +59,7 @@ public ContainerHealthState getHealthState() {
}
/**
- * Stable filter name segment used in export TAR and shard file names.
+ * Stable filter name segment used in export TAR and part file names.
*/
public String getValue() {
return value;
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/export/TestExportFileManager.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/export/TestExportFileManager.java
index 51bb88a56f5..a97255ba198 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/export/TestExportFileManager.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/export/TestExportFileManager.java
@@ -20,15 +20,25 @@
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import java.io.BufferedWriter;
import java.io.File;
+import java.io.InputStream;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.UUID;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
+import java.util.zip.GZIPInputStream;
+import org.apache.commons.compress.archivers.ArchiveInputStream;
+import org.apache.commons.compress.archivers.tar.TarArchiveEntry;
+import org.apache.commons.io.FileUtils;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleState;
import org.apache.hadoop.hdds.scm.container.ContainerHealthState;
+import org.apache.hadoop.hdds.utils.Archiver;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -38,6 +48,8 @@
*/
public class TestExportFileManager {
+ private static final String TEST_JOB_START_TIME = "2026-01-01-12-00-00";
+
@TempDir
private File tempDir;
@@ -57,12 +69,17 @@ public void testExportScopeUsesAnyForNullFilters() {
ExportScope.of(LifeCycleState.OPEN, null).getValue());
}
+ @Test
+ public void testRejectMissingFilters() {
+ assertThrows(IllegalArgumentException.class, () -> ExportScope.of(null,
null));
+ }
+
@Test
public void testResolveArchiveFile() {
ExportJob.Id jobId = ExportJob.Id.newId();
ExportScope scope = ExportScope.of(null, ContainerHealthState.MISSING);
- File archive = fileManager.resolveArchiveFile(scope, "20260101T120000Z",
jobId);
-
assertTrue(archive.getName().contains("health-MISSING_lifecycle-ANY_20260101T120000Z"));
+ File archive = fileManager.resolveArchiveFile(scope, TEST_JOB_START_TIME,
jobId);
+ assertTrue(archive.getName().contains("health-MISSING_lifecycle-ANY_" +
TEST_JOB_START_TIME));
assertTrue(archive.getName().endsWith(ExportFileManager.EXPORT_ARCHIVE_JOB_INFIX
+ jobId.getValue()
+ ExportFileManager.EXPORT_ARCHIVE_SUFFIX));
}
@@ -71,40 +88,40 @@ public void testResolveArchiveFile() {
public void testResolveArchiveTempFile() {
ExportJob.Id jobId = ExportJob.Id.newId();
ExportScope scope = ExportScope.of(null, ContainerHealthState.MISSING);
- File tempFile = fileManager.resolveArchiveTempFile(scope,
"20260101T120000Z", jobId);
+ File tempFile = fileManager.resolveArchiveTempFile(scope,
TEST_JOB_START_TIME, jobId);
assertTrue(tempFile.getName().endsWith(ExportFileManager.EXPORT_ARCHIVE_TMP_SUFFIX));
}
@Test
public void testJobIdFromArchiveFileName() {
String jobId = UUID.randomUUID().toString();
- String fileName =
"container-ids_health-MISSING_lifecycle-ANY_20260101T120000Z"
+ String fileName = "container-ids_health-MISSING_lifecycle-ANY_" +
TEST_JOB_START_TIME
+ ExportFileManager.EXPORT_ARCHIVE_JOB_INFIX + jobId +
ExportFileManager.EXPORT_ARCHIVE_SUFFIX;
assertEquals(ExportJob.Id.of(jobId),
ExportFileManager.jobIdFromArchiveFileName(fileName));
-
assertNull(ExportFileManager.jobIdFromArchiveFileName("container-ids_health-MISSING_lifecycle-ANY_20260101T120000Z"
- + ExportFileManager.EXPORT_ARCHIVE_SUFFIX));
+
assertNull(ExportFileManager.jobIdFromArchiveFileName("container-ids_health-MISSING_lifecycle-ANY_"
+ + TEST_JOB_START_TIME + ExportFileManager.EXPORT_ARCHIVE_SUFFIX));
}
@Test
- public void testArchiveTimestampFromArchiveFileName() {
- String fileName =
"container-ids_health-MISSING_lifecycle-ANY_20260101T120000Z"
+ public void testJobStartTimeFromArchiveFileName() {
+ String fileName = "container-ids_health-MISSING_lifecycle-ANY_" +
TEST_JOB_START_TIME
+ ExportFileManager.EXPORT_ARCHIVE_JOB_INFIX + UUID.randomUUID() +
ExportFileManager.EXPORT_ARCHIVE_SUFFIX;
- assertEquals("20260101T120000Z",
ExportFileManager.archiveTimestampFromArchiveFileName(fileName));
+ assertEquals(TEST_JOB_START_TIME,
ExportFileManager.jobStartTimeFromArchiveFileName(fileName));
}
@Test
public void testListCompletedArchivePaths() throws Exception {
ExportScope scope = ExportScope.of(null, ContainerHealthState.MISSING);
ExportJob.Id olderJobId = ExportJob.Id.newId();
- File olderArchive = fileManager.resolveArchiveFile(scope,
"20260101T120000Z", olderJobId);
+ File olderArchive = fileManager.resolveArchiveFile(scope,
"2026-01-01-12-00-00", olderJobId);
assertTrue(olderArchive.createNewFile());
assertTrue(olderArchive.setLastModified(2_000L));
ExportJob.Id newerJobId = ExportJob.Id.newId();
- File newerArchive = fileManager.resolveArchiveFile(scope,
"20260101T120001Z", newerJobId);
+ File newerArchive = fileManager.resolveArchiveFile(scope,
"2026-01-01-12-00-01", newerJobId);
assertTrue(newerArchive.createNewFile());
assertTrue(newerArchive.setLastModified(1_000L));
ExportJob.Id tempJobId = ExportJob.Id.newId();
- File tempArchive = fileManager.resolveArchiveTempFile(scope,
"20260101T120002Z", tempJobId);
+ File tempArchive = fileManager.resolveArchiveTempFile(scope,
"2026-01-01-12-00-02", tempJobId);
assertTrue(tempArchive.createNewFile());
List<String> completedPaths = fileManager.listCompletedArchivePaths();
@@ -113,6 +130,53 @@ public void testListCompletedArchivePaths() throws
Exception {
assertEquals(newerArchive.getAbsolutePath(), completedPaths.get(1));
}
+ @Test
+ public void testWriteArchiveFromPartFiles() throws Exception {
+ ExportJob.Id jobId = ExportJob.Id.newId();
+ ExportScope scope = ExportScope.of(null, ContainerHealthState.MISSING);
+ ExportJob job = new ExportJob(jobId, scope, TEST_JOB_START_TIME);
+ File archive = fileManager.resolveArchiveFile(scope, TEST_JOB_START_TIME,
jobId);
+
+ fileManager.createJobDirectory(jobId);
+ try (BufferedWriter writer = fileManager.newPartWriter(jobId,
job.partFileName(1))) {
+ job.writeMetadataHeader(writer, 1, 1L);
+ writer.write("1\n2\n");
+ }
+
+ fileManager.writeArchive(jobId, archive.getAbsolutePath());
+ fileManager.deleteJobDirectory(jobId);
+
+ assertTrue(archive.exists());
+ Path extractDir = Files.createTempDirectory("export-archive");
+ try {
+ extractGzTar(archive, extractDir);
+ List<String> partNames;
+ try (Stream<Path> stream = Files.list(extractDir)) {
+ partNames = stream.map(path ->
path.getFileName().toString()).collect(Collectors.toList());
+ }
+ assertEquals(1, partNames.size());
+ assertTrue(partNames.get(0).endsWith("part001.txt"));
+ } finally {
+ FileUtils.deleteQuietly(extractDir.toFile());
+ }
+ }
+
+ @Test
+ public void testCleanupFailedJobRemovesJobDirAndTempArchive() throws
Exception {
+ ExportJob.Id jobId = ExportJob.Id.newId();
+ ExportScope scope = ExportScope.of(null, ContainerHealthState.MISSING);
+ File archive = fileManager.resolveArchiveFile(scope, TEST_JOB_START_TIME,
jobId);
+ File tempArchive = fileManager.resolveArchiveTempFile(scope,
TEST_JOB_START_TIME, jobId);
+ assertTrue(tempArchive.createNewFile());
+
+ fileManager.createJobDirectory(jobId);
+ fileManager.cleanupFailedJob(jobId, archive.getAbsolutePath());
+
+
assertFalse(Files.exists(tempDir.toPath().resolve(ExportFileManager.exportJobDirName(jobId))));
+ assertFalse(tempArchive.exists());
+ assertFalse(archive.exists());
+ }
+
@Test
public void testOrphanJobDirRemovedOnStartup() throws Exception {
ExportJob.Id jobId = ExportJob.Id.newId();
@@ -130,7 +194,7 @@ public void testIncompleteExportArtifactsRemovedOnStartup()
throws Exception {
Path jobDir =
tempDir.toPath().resolve(ExportFileManager.exportJobDirName(jobId));
Files.createDirectories(jobDir);
ExportScope scope = ExportScope.of(null, ContainerHealthState.MISSING);
- File partialArchiveTemp = fileManager.resolveArchiveTempFile(scope,
"20260101T000000Z", jobId);
+ File partialArchiveTemp = fileManager.resolveArchiveTempFile(scope,
"2026-01-01-00-00-00", jobId);
assertTrue(partialArchiveTemp.createNewFile());
fileManager.start();
@@ -143,7 +207,7 @@ public void testIncompleteExportArtifactsRemovedOnStartup()
throws Exception {
public void testOrphanJobDirDoesNotDeleteCompletedTar() throws Exception {
ExportJob.Id jobId = ExportJob.Id.newId();
ExportScope scope = ExportScope.of(null, ContainerHealthState.MISSING);
- File completedArchive = fileManager.resolveArchiveFile(scope,
"20260101T000000Z", jobId);
+ File completedArchive = fileManager.resolveArchiveFile(scope,
"2026-01-01-00-00-00", jobId);
assertTrue(completedArchive.createNewFile());
Path orphanJobDir =
tempDir.toPath().resolve(ExportFileManager.exportJobDirName(jobId));
Files.createDirectories(orphanJobDir);
@@ -153,4 +217,15 @@ public void testOrphanJobDirDoesNotDeleteCompletedTar()
throws Exception {
assertTrue(completedArchive.exists());
assertFalse(Files.exists(orphanJobDir));
}
+
+ private static void extractGzTar(File archive, Path extractDir) throws
Exception {
+ Files.createDirectories(extractDir);
+ try (InputStream in = new
GZIPInputStream(Files.newInputStream(archive.toPath()));
+ ArchiveInputStream<TarArchiveEntry> tarIn = Archiver.untar(in)) {
+ TarArchiveEntry entry;
+ while ((entry = tarIn.getNextEntry()) != null) {
+ Archiver.extractEntry(entry, tarIn, entry.getSize(), extractDir,
extractDir.resolve(entry.getName()));
+ }
+ }
+ }
}
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/export/TestExportJob.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/export/TestExportJob.java
new file mode 100644
index 00000000000..fdad57b4dc3
--- /dev/null
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/export/TestExportJob.java
@@ -0,0 +1,67 @@
+/*
+ * 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.hdds.scm.container.export;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.BufferedWriter;
+import java.io.File;
+import java.nio.file.Files;
+import java.util.List;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleState;
+import org.apache.hadoop.hdds.scm.container.ContainerHealthState;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+/**
+ * Tests for {@link ExportJob}.
+ */
+public class TestExportJob {
+
+ private static final String TEST_JOB_START_TIME = "2026-01-01-12-00-00";
+
+ @TempDir
+ private File tempDir;
+
+ @Test
+ public void testPartFileName() {
+ ExportJob job = newJob(ContainerHealthState.MISSING, null);
+
assertEquals("container-ids_health-MISSING_lifecycle-ANY_2026-01-01-12-00-00_part001.txt",
job.partFileName(1));
+ }
+
+ @Test
+ public void testWriteMetadataHeader() throws Exception {
+ ExportJob job = newJob(ContainerHealthState.MISSING, LifeCycleState.OPEN);
+ File headerFile = new File(tempDir, "part-header.txt");
+ try (BufferedWriter writer = Files.newBufferedWriter(headerFile.toPath()))
{
+ job.writeMetadataHeader(writer, 2, 42L);
+ }
+ List<String> lines = Files.readAllLines(headerFile.toPath());
+ assertTrue(lines.contains("# jobStartTime=" + TEST_JOB_START_TIME));
+ assertTrue(lines.contains("# healthState=MISSING"));
+ assertTrue(lines.contains("# lifecycleState=OPEN"));
+ assertTrue(lines.contains("# startContainerId=42"));
+ assertTrue(lines.contains("# part=2"));
+ }
+
+ private static ExportJob newJob(ContainerHealthState healthState,
LifeCycleState lifeCycleState) {
+ ExportScope scope = ExportScope.of(lifeCycleState, healthState);
+ return new ExportJob(ExportJob.Id.newId(), scope, TEST_JOB_START_TIME);
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]