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]

Reply via email to