From 8ade6a7ac3f98de2e323120724e9e71b4180752f Mon Sep 17 00:00:00 2001 From: Bruno Volpato Date: Thu, 2 Apr 2026 10:24:32 -0400 Subject: [PATCH 1/2] Fix flaky FileIOTest.testMatchWatchForNewFiles test under CI filesystems addresses #19480 The `FileIOTest.testMatchWatchForNewFiles` test occasionally flakes in the CI environment because the `updOptions` configuration in `CopyFilesFn` does not preserve file attributes when overwriting existing files. In some CI filesystems, this causes the copied file to register a `lastModifiedMillis` timestamp of `0`. Thus, when `ExtractFilenameAndLastUpdateFn` parses this file, it throws a `RuntimeException` at `FileIO.java:800`, failing the pipeline run. This PR adds `StandardCopyOption.COPY_ATTRIBUTES` to preserve the file's original timestamps, avoiding the exception. --- .../core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java index dffc4943bfab..dc4457bb4a20 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java @@ -249,7 +249,9 @@ public void processElement(ProcessContext context, @StateId("count") ValueState< context.output(Objects.requireNonNull(context.element()).getValue()); CopyOption[] cpOptions = {StandardCopyOption.COPY_ATTRIBUTES}; - CopyOption[] updOptions = {StandardCopyOption.REPLACE_EXISTING}; + CopyOption[] updOptions = { + StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.COPY_ATTRIBUTES + }; final Path sourcePath = Paths.get(sourcePathStr); final Path watchPath = Paths.get(watchPathStr); From 2d8a2bb576cff92602e6a5dc63aa3a8492838c10 Mon Sep 17 00:00:00 2001 From: Bruno Volpato Date: Tue, 5 May 2026 01:47:28 -0400 Subject: [PATCH 2/2] Stabilize FileIOTest updated-file timestamp assertions --- .../org/apache/beam/sdk/io/FileIOTest.java | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 deletions(-) diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java index dc4457bb4a20..4f7c0eb62ff9 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileIOTest.java @@ -232,9 +232,10 @@ public void testMatchAllDisallowEmptyNonWildcard() throws IOException { /** DoFn that copy test files from source to watch path. */ private static class CopyFilesFn extends DoFn, MatchResult.Metadata> { - public CopyFilesFn(Path sourcePath, Path watchPath) { + public CopyFilesFn(Path sourcePath, Path watchPath, long baseTimestampMillis) { this.sourcePathStr = sourcePath.toString(); this.watchPathStr = watchPath.toString(); + this.baseTimestampMillis = baseTimestampMillis; } @StateId("count") @@ -257,10 +258,16 @@ public void processElement(ProcessContext context, @StateId("count") ValueState< if (0 == current) { Thread.sleep(100); + // Ensure overwrite updates get a distinct mtime even when COPY_ATTRIBUTES is enabled. + Files.setLastModifiedTime( + sourcePath.resolve("first"), FileTime.fromMillis(baseTimestampMillis + 2000)); Files.copy(sourcePath.resolve("first"), watchPath.resolve("first"), updOptions); Files.copy(sourcePath.resolve("second"), watchPath.resolve("second"), cpOptions); } else if (1 == current) { Thread.sleep(100); + FileTime updateTime = FileTime.fromMillis(baseTimestampMillis + 4000); + Files.setLastModifiedTime(sourcePath.resolve("first"), updateTime); + Files.setLastModifiedTime(sourcePath.resolve("second"), updateTime); Files.copy(sourcePath.resolve("first"), watchPath.resolve("first"), updOptions); Files.copy(sourcePath.resolve("second"), watchPath.resolve("second"), updOptions); Files.copy(sourcePath.resolve("third"), watchPath.resolve("third"), cpOptions); @@ -271,6 +278,7 @@ public void processElement(ProcessContext context, @StateId("count") ValueState< // Member variables need to be serializable. private final String sourcePathStr; private final String watchPathStr; + private final long baseTimestampMillis; } @Test @@ -282,6 +290,12 @@ public void testMatchWatchForNewFiles() throws IOException, InterruptedException Files.write(sourcePath.resolve("first"), new byte[42]); Files.write(sourcePath.resolve("second"), new byte[37]); Files.write(sourcePath.resolve("third"), new byte[99]); + // Keep controlled mtimes in the past so updates are distinct without future timestamps. + long baseTimestampMillis = System.currentTimeMillis() - Duration.standardMinutes(1).getMillis(); + FileTime baseTimestamp = FileTime.fromMillis(baseTimestampMillis); + Files.setLastModifiedTime(sourcePath.resolve("first"), baseTimestamp); + Files.setLastModifiedTime(sourcePath.resolve("second"), baseTimestamp); + Files.setLastModifiedTime(sourcePath.resolve("third"), baseTimestamp); // Create a "watch" directory that the pipeline will copy files into. final Path watchPath = tmpFolder.getRoot().toPath().resolve("watch"); @@ -336,7 +350,7 @@ public void testMatchWatchForNewFiles() throws IOException, InterruptedException TypeDescriptors.strings(), TypeDescriptor.of(MatchResult.Metadata.class))) .via((metadata) -> KV.of("dumb key", metadata))) - .apply(ParDo.of(new CopyFilesFn(sourcePath, watchPath))); + .apply(ParDo.of(new CopyFilesFn(sourcePath, watchPath, baseTimestampMillis))); assertEquals(PCollection.IsBounded.UNBOUNDED, matchMetadata.isBounded()); assertEquals(PCollection.IsBounded.UNBOUNDED, matchAllMetadata.isBounded());