This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11781-46b75fe8b92097b1d36c75e2df9e82f9d0ae73c1 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit e896b9cbe4cce59b38ec76969aaa7a0229c54869 Author: Jast <[email protected]> AuthorDate: Mon Sep 28 17:39:22 2026 +0000 [Fix][Zeta] Restore legacy job status WAL entries (#11781) Co-authored-by: zhangshenghang <[email protected]> Co-authored-by: zhangshenghang <[email protected]> Co-authored-by: zhangshenghang <[email protected]> --- .../imap-storage-plugins/imap-storage-file/pom.xml | 6 ++++ .../engine/imap/storage/file/common/WALReader.java | 40 +++++++++++++++++++++- .../file/common/WALReaderAndWriterTest.java | 26 ++++++++++++++ 3 files changed, 71 insertions(+), 1 deletion(-) diff --git a/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/pom.xml b/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/pom.xml index 96069b75c0..de3296f3de 100644 --- a/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/pom.xml +++ b/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/pom.xml @@ -43,6 +43,12 @@ <artifactId>serializer-protobuf</artifactId> <version>${project.version}</version> </dependency> + <dependency> + <groupId>org.apache.seatunnel</groupId> + <artifactId>seatunnel-engine-common</artifactId> + <version>${project.version}</version> + <scope>test</scope> + </dependency> <!-- hadoop jar --> <dependency> <groupId>org.apache.seatunnel</groupId> diff --git a/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/src/main/java/org/apache/seatunnel/engine/imap/storage/file/common/WALReader.java b/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/src/main/java/org/apache/seatunnel/engine/imap/storage/file/common/WALReader.java index f026e2f940..414d04b25f 100644 --- a/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/src/main/java/org/apache/seatunnel/engine/imap/storage/file/common/WALReader.java +++ b/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/src/main/java/org/apache/seatunnel/engine/imap/storage/file/common/WALReader.java @@ -42,6 +42,33 @@ import java.util.Set; import static org.apache.seatunnel.engine.imap.storage.file.common.LatestMutationAccumulator.SerializedKey; public class WALReader { + /** + * Maps the fully-qualified class names recorded in WAL entries written by SeaTunnel versions + * older than #9689 (which moved the job status model from {@code seatunnel-engine-core} to + * {@code seatunnel-engine-common}) to their current class names. + * + * <p>Without this mapping, restarting a cluster that persists job state through the file-based + * IMap storage fails with {@code ClassNotFoundException} while replaying existing WAL entries, + * see #9928. Removing an entry from this table breaks upgrades from any release that still + * wrote the old class name, so keep it in place and extend it whenever a persisted job model + * class is moved again. + */ + private static final Map<String, String> LEGACY_CLASS_NAME_MAPPINGS; + + static { + Map<String, String> legacyClassNameMappings = new HashMap<>(); + legacyClassNameMappings.put( + "org.apache.seatunnel.engine.core.job.JobResult", + "org.apache.seatunnel.engine.common.job.JobResult"); + legacyClassNameMappings.put( + "org.apache.seatunnel.engine.core.job.JobStatus", + "org.apache.seatunnel.engine.common.job.JobStatus"); + legacyClassNameMappings.put( + "org.apache.seatunnel.engine.core.job.JobStatusData", + "org.apache.seatunnel.engine.common.job.JobStatusData"); + LEGACY_CLASS_NAME_MAPPINGS = Collections.unmodifiableMap(legacyClassNameMappings); + } + private final Serializer serializer; private final IFileReader<IMapFileData> fileReader; @@ -108,7 +135,7 @@ public class WALReader { private Object deserializeData(byte[] data, String className) { try { - Class<?> clazz = ClassUtils.getClass(className); + Class<?> clazz = ClassUtils.getClass(resolveClassName(className)); try { return serializer.deserialize(data, clazz); } catch (IOException e) { @@ -123,4 +150,15 @@ public class WALReader { e, "deserialize data error, class name is {}", className); } } + + /** + * Resolves class names written by SeaTunnel versions before the job model moved to + * engine-common. + * + * @param className class name recorded in the WAL + * @return class name available in the current distribution + */ + private static String resolveClassName(String className) { + return LEGACY_CLASS_NAME_MAPPINGS.getOrDefault(className, className); + } } diff --git a/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/src/test/java/org/apache/seatunnel/engine/imap/storage/file/common/WALReaderAndWriterTest.java b/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/src/test/java/org/apache/seatunnel/engine/imap/storage/file/common/WALReaderAndWriterTest.java index 6b1fb4b7ff..3b70327963 100644 --- a/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/src/test/java/org/apache/seatunnel/engine/imap/storage/file/common/WALReaderAndWriterTest.java +++ b/seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file/src/test/java/org/apache/seatunnel/engine/imap/storage/file/common/WALReaderAndWriterTest.java @@ -20,6 +20,7 @@ package org.apache.seatunnel.engine.imap.storage.file.common; +import org.apache.seatunnel.engine.common.job.JobStatus; import org.apache.seatunnel.engine.imap.storage.file.bean.IMapFileData; import org.apache.seatunnel.engine.imap.storage.file.config.FileConfiguration; import org.apache.seatunnel.engine.serializer.api.Serializer; @@ -55,6 +56,7 @@ public class WALReaderAndWriterTest { private static final Path PARENT_PATH = new Path("/tmp/9/"); private static final Path SAME_TIMESTAMP_TOMBSTONE_PATH = new Path("/tmp/imap-wal-same-timestamp-tombstone/"); + private static final Path LEGACY_JOB_STATUS_PATH = new Path("/tmp/imap-wal-legacy-job-status/"); private static final Serializer SERIALIZER = new ProtoStuffSerializer(); @BeforeAll @@ -240,10 +242,34 @@ public class WALReaderAndWriterTest { Assertions.assertFalse(reader.loadAllKeys(SAME_TIMESTAMP_TOMBSTONE_PATH).contains(key)); } + @Test + public void testReaderLoadsJobStatusWrittenWithLegacyClassName() throws Exception { + String key = "job-status"; + try (WALWriter writer = + new WALWriter(FS, FileConfiguration.HDFS, LEGACY_JOB_STATUS_PATH, SERIALIZER)) { + writer.write( + IMapFileData.builder() + .key(SERIALIZER.serialize(key)) + .keyClassName(String.class.getName()) + .value(SERIALIZER.serialize(JobStatus.RUNNING)) + .valueClassName("org.apache.seatunnel.engine.core.job.JobStatus") + .deleted(false) + .timestamp(System.nanoTime()) + .build()); + } + + WALReader reader = new WALReader(FS, FileConfiguration.HDFS, SERIALIZER); + Map<Object, Object> data = reader.loadAllData(LEGACY_JOB_STATUS_PATH, new HashSet<>()); + + Assertions.assertEquals(JobStatus.class, data.get(key).getClass()); + Assertions.assertEquals(JobStatus.RUNNING.name(), data.get(key).toString()); + } + @AfterAll public static void close() throws IOException { FS.delete(PARENT_PATH, true); FS.delete(SAME_TIMESTAMP_TOMBSTONE_PATH, true); + FS.delete(LEGACY_JOB_STATUS_PATH, true); FS.close(); } }
