From 88e2f20fe72c024cc881396c2ac0707159d5591d Mon Sep 17 00:00:00 2001 From: Andrew Kabas Date: Thu, 8 Jan 2026 17:47:12 -0500 Subject: [PATCH 1/8] Report source lineage from HadoopFormatIO --- .../sdk/io/hadoop/format/HadoopFormatIO.java | 42 +++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java index 0412e4286bb8..eed4419eca36 100644 --- a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java +++ b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java @@ -52,6 +52,9 @@ import org.apache.beam.sdk.coders.CoderRegistry; import org.apache.beam.sdk.coders.KvCoder; import org.apache.beam.sdk.io.BoundedSource; +import org.apache.beam.sdk.io.FileSystem; +import org.apache.beam.sdk.io.FileSystems; +import org.apache.beam.sdk.io.fs.ResourceId; import org.apache.beam.sdk.io.hadoop.SerializableConfiguration; import org.apache.beam.sdk.io.hadoop.WritableCoder; import org.apache.beam.sdk.options.PipelineOptions; @@ -83,6 +86,7 @@ import org.apache.hadoop.io.ObjectWritable; import org.apache.hadoop.io.Writable; import org.apache.hadoop.mapred.FileAlreadyExistsException; +import org.apache.hadoop.mapred.FileSplit; import org.apache.hadoop.mapreduce.InputFormat; import org.apache.hadoop.mapreduce.InputSplit; import org.apache.hadoop.mapreduce.Job; @@ -462,6 +466,7 @@ public Read withConfiguration(Configuration configuration) { if (getValueTranslationFunction() == null) { builder.setValueTypeDescriptor((TypeDescriptor) inputFormatValueClass); } + return builder.build(); } @@ -724,7 +729,11 @@ public List>> split(long desiredBundleSizeBytes, Pipeline LOG.info("Not splitting source {} because source is already split.", this); return ImmutableList.of(this); } + computeSplitsIfNecessary(); + + reportSourceLineage(inputSplits); + LOG.info( "Generated {} splits. Size of first split is {} ", inputSplits.size(), @@ -744,6 +753,39 @@ public List>> split(long desiredBundleSizeBytes, Pipeline .collect(Collectors.toList()); } + private void reportSourceLineage(final List inputSplits) { + List fileResources = new ArrayList<>(); + + for (SerializableSplit split : inputSplits) { + InputSplit inputSplit = split.getSplit(); + + if (inputSplit instanceof FileSplit) { + String pathString = ((FileSplit) inputSplit).getPath().toString(); + ResourceId resourceId = FileSystems.matchNewResource(pathString, false); + fileResources.add(resourceId); + } + } + + if (fileResources.size() <= 100) { + for (ResourceId resource : fileResources) { + FileSystems.reportSourceLineage(resource); + } + } else { + HashSet uniqueDirs = new HashSet<>(); + for (ResourceId resource : fileResources) { + ResourceId dir = resource.getCurrentDirectory(); + uniqueDirs.add(dir); + if (uniqueDirs.size() > 100) { + FileSystems.reportSourceLineage(dir, FileSystem.LineageLevel.TOP_LEVEL); + return; + } + } + for (ResourceId dir : uniqueDirs) { + FileSystems.reportSourceLineage(dir); + } + } + } + @Override public long getEstimatedSizeBytes(PipelineOptions po) throws Exception { if (inputSplit == null) { From 4bf30f62e621d2687cb15c1f39396ea5b90a55a2 Mon Sep 17 00:00:00 2001 From: Andrew Kabas Date: Thu, 8 Jan 2026 18:08:04 -0500 Subject: [PATCH 2/8] Extract shared code into FileSystems public API --- .../apache/beam/sdk/io/FileBasedSource.java | 34 +++---------------- .../org/apache/beam/sdk/io/FileSystems.java | 32 +++++++++++++++++ .../sdk/io/hadoop/format/HadoopFormatIO.java | 21 ++---------- 3 files changed, 38 insertions(+), 49 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSource.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSource.java index 8d6e52c64a52..4d65fbbd5e93 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSource.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSource.java @@ -26,12 +26,11 @@ import java.nio.channels.ReadableByteChannel; import java.nio.channels.SeekableByteChannel; import java.util.ArrayList; -import java.util.HashSet; import java.util.List; import java.util.ListIterator; import java.util.NoSuchElementException; import java.util.concurrent.atomic.AtomicReference; -import org.apache.beam.sdk.io.FileSystem.LineageLevel; +import java.util.stream.Collectors; import org.apache.beam.sdk.io.fs.EmptyMatchTreatment; import org.apache.beam.sdk.io.fs.MatchResult; import org.apache.beam.sdk.io.fs.MatchResult.Metadata; @@ -318,35 +317,10 @@ public final List> split( } } - /** - * Report source Lineage. Due to the size limit of Beam metrics, report full file name or only dir - * depend on the number of files. - * - *

- Number of files<=100, report full file paths; - * - *

- Number of directory<=100, report directory names (one level up); - * - *

- Otherwise, report top level only. - */ private static void reportSourceLineage(List expandedFiles) { - if (expandedFiles.size() <= 100) { - for (Metadata metadata : expandedFiles) { - FileSystems.reportSourceLineage(metadata.resourceId()); - } - } else { - HashSet uniqueDirs = new HashSet<>(); - for (Metadata metadata : expandedFiles) { - ResourceId dir = metadata.resourceId().getCurrentDirectory(); - uniqueDirs.add(dir); - if (uniqueDirs.size() > 100) { - FileSystems.reportSourceLineage(dir, LineageLevel.TOP_LEVEL); - return; - } - } - for (ResourceId uniqueDir : uniqueDirs) { - FileSystems.reportSourceLineage(uniqueDir); - } - } + List resourceIds = + expandedFiles.stream().map(Metadata::resourceId).collect(Collectors.toList()); + FileSystems.reportSourceLineage(resourceIds); } /** diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java index 6133ca9fdb39..4a79473ee03d 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java @@ -29,6 +29,7 @@ import java.util.ArrayList; import java.util.Collection; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Map.Entry; @@ -398,6 +399,37 @@ public ResourceId apply(@Nonnull Metadata input) { .delete(resourceIdsToDelete); } + /** + * Report source {@link Lineage} metrics for multiple resource ids. Due to the size limit of Beam + * metrics, report full file name or only dir depend on the number of files. + * + *

- Number of files<=100, report full file paths; + * + *

- Number of directory<=100, report directory names (one level up); + * + *

- Otherwise, report top level only. + */ + public static void reportSourceLineage(List resourceIds) { + if (resourceIds.size() <= 100) { + for (ResourceId resourceId : resourceIds) { + FileSystems.reportSourceLineage(resourceId); + } + } else { + HashSet uniqueDirs = new HashSet<>(); + for (ResourceId resourceId : resourceIds) { + ResourceId dir = resourceId.getCurrentDirectory(); + uniqueDirs.add(dir); + if (uniqueDirs.size() > 100) { + FileSystems.reportSourceLineage(dir, LineageLevel.TOP_LEVEL); + return; + } + } + for (ResourceId uniqueDir : uniqueDirs) { + FileSystems.reportSourceLineage(uniqueDir); + } + } + } + /** Report source {@link Lineage} metrics for resource id. */ public static void reportSourceLineage(ResourceId resourceId) { reportSourceLineage(resourceId, LineageLevel.FILE); diff --git a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java index eed4419eca36..cecf55ca99f3 100644 --- a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java +++ b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java @@ -52,7 +52,6 @@ import org.apache.beam.sdk.coders.CoderRegistry; import org.apache.beam.sdk.coders.KvCoder; import org.apache.beam.sdk.io.BoundedSource; -import org.apache.beam.sdk.io.FileSystem; import org.apache.beam.sdk.io.FileSystems; import org.apache.beam.sdk.io.fs.ResourceId; import org.apache.beam.sdk.io.hadoop.SerializableConfiguration; @@ -753,6 +752,7 @@ public List>> split(long desiredBundleSizeBytes, Pipeline .collect(Collectors.toList()); } + /** Report only file-based sources */ private void reportSourceLineage(final List inputSplits) { List fileResources = new ArrayList<>(); @@ -766,24 +766,7 @@ private void reportSourceLineage(final List inputSplits) { } } - if (fileResources.size() <= 100) { - for (ResourceId resource : fileResources) { - FileSystems.reportSourceLineage(resource); - } - } else { - HashSet uniqueDirs = new HashSet<>(); - for (ResourceId resource : fileResources) { - ResourceId dir = resource.getCurrentDirectory(); - uniqueDirs.add(dir); - if (uniqueDirs.size() > 100) { - FileSystems.reportSourceLineage(dir, FileSystem.LineageLevel.TOP_LEVEL); - return; - } - } - for (ResourceId dir : uniqueDirs) { - FileSystems.reportSourceLineage(dir); - } - } + FileSystems.reportSourceLineage(fileResources); } @Override From 90b6cabd155b0c73275ced87a370b978f5102784 Mon Sep 17 00:00:00 2001 From: Andrew Kabas Date: Fri, 9 Jan 2026 16:37:07 -0500 Subject: [PATCH 3/8] A little refactoring --- .../org/apache/beam/sdk/io/FileSystems.java | 5 +++-- .../sdk/io/hadoop/format/HadoopFormatIO.java | 22 ++++++------------- 2 files changed, 10 insertions(+), 17 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java index 4a79473ee03d..eb401ea5fd53 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java @@ -410,7 +410,8 @@ public ResourceId apply(@Nonnull Metadata input) { *

- Otherwise, report top level only. */ public static void reportSourceLineage(List resourceIds) { - if (resourceIds.size() <= 100) { + final int MAX_LINEAGE_TARGETS = 100; + if (resourceIds.size() <= MAX_LINEAGE_TARGETS) { for (ResourceId resourceId : resourceIds) { FileSystems.reportSourceLineage(resourceId); } @@ -419,7 +420,7 @@ public static void reportSourceLineage(List resourceIds) { for (ResourceId resourceId : resourceIds) { ResourceId dir = resourceId.getCurrentDirectory(); uniqueDirs.add(dir); - if (uniqueDirs.size() > 100) { + if (uniqueDirs.size() > MAX_LINEAGE_TARGETS) { FileSystems.reportSourceLineage(dir, LineageLevel.TOP_LEVEL); return; } diff --git a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java index cecf55ca99f3..0b105026773d 100644 --- a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java +++ b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java @@ -465,7 +465,6 @@ public Read withConfiguration(Configuration configuration) { if (getValueTranslationFunction() == null) { builder.setValueTypeDescriptor((TypeDescriptor) inputFormatValueClass); } - return builder.build(); } @@ -728,11 +727,8 @@ public List>> split(long desiredBundleSizeBytes, Pipeline LOG.info("Not splitting source {} because source is already split.", this); return ImmutableList.of(this); } - computeSplitsIfNecessary(); - reportSourceLineage(inputSplits); - LOG.info( "Generated {} splits. Size of first split is {} ", inputSplits.size(), @@ -754,17 +750,13 @@ public List>> split(long desiredBundleSizeBytes, Pipeline /** Report only file-based sources */ private void reportSourceLineage(final List inputSplits) { - List fileResources = new ArrayList<>(); - - for (SerializableSplit split : inputSplits) { - InputSplit inputSplit = split.getSplit(); - - if (inputSplit instanceof FileSplit) { - String pathString = ((FileSplit) inputSplit).getPath().toString(); - ResourceId resourceId = FileSystems.matchNewResource(pathString, false); - fileResources.add(resourceId); - } - } + List fileResources = + inputSplits.stream() + .map(SerializableSplit::getSplit) + .filter(FileSplit.class::isInstance) + .map(FileSplit.class::cast) + .map(fileSplit -> FileSystems.matchNewResource(fileSplit.getPath().toString(), false)) + .collect(Collectors.toList()); FileSystems.reportSourceLineage(fileResources); } From 5108e085b93fecf8d1cc8795e249a1172d45ac2e Mon Sep 17 00:00:00 2001 From: Andrew Kabas Date: Fri, 9 Jan 2026 19:43:05 -0500 Subject: [PATCH 4/8] Fix import in HadoopFormatIO --- .../org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java index 0b105026773d..50a53860db7d 100644 --- a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java +++ b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java @@ -85,7 +85,6 @@ import org.apache.hadoop.io.ObjectWritable; import org.apache.hadoop.io.Writable; import org.apache.hadoop.mapred.FileAlreadyExistsException; -import org.apache.hadoop.mapred.FileSplit; import org.apache.hadoop.mapreduce.InputFormat; import org.apache.hadoop.mapreduce.InputSplit; import org.apache.hadoop.mapreduce.Job; @@ -100,6 +99,7 @@ import org.apache.hadoop.mapreduce.TaskAttemptID; import org.apache.hadoop.mapreduce.TaskID; import org.apache.hadoop.mapreduce.lib.db.DBConfiguration; +import org.apache.hadoop.mapreduce.lib.input.FileSplit; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import org.apache.hadoop.mapreduce.task.JobContextImpl; import org.apache.hadoop.mapreduce.task.TaskAttemptContextImpl; From 15b8ae06dbe4119fb0686d654d3f6de12f0dbce2 Mon Sep 17 00:00:00 2001 From: Andrew Kabas Date: Mon, 12 Jan 2026 11:25:38 -0500 Subject: [PATCH 5/8] check style --- .../src/main/java/org/apache/beam/sdk/io/FileSystems.java | 6 +++--- .../apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java index eb401ea5fd53..e93f4efa8e34 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java @@ -410,8 +410,8 @@ public ResourceId apply(@Nonnull Metadata input) { *

- Otherwise, report top level only. */ public static void reportSourceLineage(List resourceIds) { - final int MAX_LINEAGE_TARGETS = 100; - if (resourceIds.size() <= MAX_LINEAGE_TARGETS) { + final int maxLineageTargets = 100; + if (resourceIds.size() <= maxLineageTargets) { for (ResourceId resourceId : resourceIds) { FileSystems.reportSourceLineage(resourceId); } @@ -420,7 +420,7 @@ public static void reportSourceLineage(List resourceIds) { for (ResourceId resourceId : resourceIds) { ResourceId dir = resourceId.getCurrentDirectory(); uniqueDirs.add(dir); - if (uniqueDirs.size() > MAX_LINEAGE_TARGETS) { + if (uniqueDirs.size() > maxLineageTargets) { FileSystems.reportSourceLineage(dir, LineageLevel.TOP_LEVEL); return; } diff --git a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java index 50a53860db7d..58780a9eb66c 100644 --- a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java +++ b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java @@ -748,7 +748,7 @@ public List>> split(long desiredBundleSizeBytes, Pipeline .collect(Collectors.toList()); } - /** Report only file-based sources */ + /** Report only file-based sources. */ private void reportSourceLineage(final List inputSplits) { List fileResources = inputSplits.stream() From b4549cad1a538832fb5b911b833dabd6c0665041 Mon Sep 17 00:00:00 2001 From: Andrew Kabas Date: Mon, 30 Mar 2026 18:09:58 -0400 Subject: [PATCH 6/8] Implement lineage in HDFS and add tests --- .../apache/beam/sdk/io/FileSystemsTest.java | 98 +++++++++++++++++++ .../beam/sdk/io/hdfs/HadoopFileSystem.java | 22 +++++ .../sdk/io/hdfs/HadoopFileSystemTest.java | 19 ++++ 3 files changed, 139 insertions(+) diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileSystemsTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileSystemsTest.java index 34567309c7d0..bd7cb413739b 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileSystemsTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileSystemsTest.java @@ -23,8 +23,11 @@ import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; import static org.junit.Assume.assumeFalse; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -35,7 +38,9 @@ import java.nio.file.NoSuchFileException; import java.nio.file.Path; import java.nio.file.Paths; +import java.util.ArrayList; import java.util.List; +import org.apache.beam.sdk.io.FileSystem.LineageLevel; import org.apache.beam.sdk.io.fs.CreateOptions; import org.apache.beam.sdk.io.fs.MatchResult; import org.apache.beam.sdk.io.fs.MoveOptions; @@ -54,6 +59,8 @@ import org.junit.rules.TemporaryFolder; import org.junit.runner.RunWith; import org.junit.runners.JUnit4; +import org.mockito.MockedStatic; +import org.mockito.Mockito; /** Tests for {@link FileSystems}. */ @RunWith(JUnit4.class) @@ -337,6 +344,97 @@ public void testMatchNewDirectory() { } } + @Test + public void testReportSourceLineageFewFiles() { + List resources = new ArrayList<>(); + for (int i = 0; i < 5; i++) { + resources.add(mockResourceId("/dir/file" + i, "/dir/")); + } + + try (MockedStatic mocked = + Mockito.mockStatic(FileSystems.class, Mockito.CALLS_REAL_METHODS)) { + mocked + .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class))) + .thenAnswer(inv -> null); + mocked + .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class), any())) + .thenAnswer(inv -> null); + + FileSystems.reportSourceLineage(resources); + + for (ResourceId r : resources) { + mocked.verify(() -> FileSystems.reportSourceLineage(r)); + } + } + } + + @Test + public void testReportSourceLineageManyFilesFewDirs() { + ResourceId dir = mockResourceId("/dir/", null); + List resources = new ArrayList<>(); + for (int i = 0; i < 150; i++) { + resources.add(mockResourceId("/dir/file" + i, dir)); + } + + try (MockedStatic mocked = + Mockito.mockStatic(FileSystems.class, Mockito.CALLS_REAL_METHODS)) { + mocked + .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class))) + .thenAnswer(inv -> null); + mocked + .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class), any())) + .thenAnswer(inv -> null); + + FileSystems.reportSourceLineage(resources); + + // Should report the unique directory, not individual files + mocked.verify(() -> FileSystems.reportSourceLineage(dir)); + mocked.verify( + () -> FileSystems.reportSourceLineage(any(ResourceId.class), any(LineageLevel.class)), + never()); + } + } + + @Test + public void testReportSourceLineageManyFilesManyDirs() { + List resources = new ArrayList<>(); + for (int i = 0; i < 150; i++) { + ResourceId dir = mockResourceId("/dir" + i + "/", null); + resources.add(mockResourceId("/dir" + i + "/file", dir)); + } + + try (MockedStatic mocked = + Mockito.mockStatic(FileSystems.class, Mockito.CALLS_REAL_METHODS)) { + mocked + .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class))) + .thenAnswer(inv -> null); + mocked + .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class), any())) + .thenAnswer(inv -> null); + + FileSystems.reportSourceLineage(resources); + + // Should fall back to TOP_LEVEL reporting + mocked.verify( + () -> FileSystems.reportSourceLineage(any(ResourceId.class), eq(LineageLevel.TOP_LEVEL))); + // Should not report individual files or directories at FILE level + mocked.verify(() -> FileSystems.reportSourceLineage(any(ResourceId.class)), never()); + } + } + + private ResourceId mockResourceId(String path, Object dir) { + ResourceId resourceId = mock(ResourceId.class); + when(resourceId.toString()).thenReturn(path); + if (dir instanceof String) { + ResourceId dirId = mock(ResourceId.class); + when(dirId.toString()).thenReturn((String) dir); + when(resourceId.getCurrentDirectory()).thenReturn(dirId); + } else if (dir instanceof ResourceId) { + when(resourceId.getCurrentDirectory()).thenReturn((ResourceId) dir); + } + return resourceId; + } + private static List toResourceIds(List paths, final boolean isDirectory) { return FluentIterable.from(paths) .transform(path -> (ResourceId) LocalResourceId.fromPath(path, isDirectory)) diff --git a/sdks/java/io/hadoop-file-system/src/main/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java b/sdks/java/io/hadoop-file-system/src/main/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java index 4de1d539e76a..e0a7793a243e 100644 --- a/sdks/java/io/hadoop-file-system/src/main/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java +++ b/sdks/java/io/hadoop-file-system/src/main/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java @@ -38,6 +38,7 @@ import org.apache.beam.sdk.io.fs.MatchResult.Metadata; import org.apache.beam.sdk.io.fs.MatchResult.Status; import org.apache.beam.sdk.io.fs.MoveOptions; +import org.apache.beam.sdk.metrics.Lineage; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; import org.apache.hadoop.conf.Configuration; @@ -336,6 +337,27 @@ protected String getScheme() { return scheme; } + @Override + protected void reportLineage(HadoopResourceId resourceId, Lineage lineage) { + reportLineage(resourceId, lineage, LineageLevel.FILE); + } + + @Override + protected void reportLineage(HadoopResourceId resourceId, Lineage lineage, LineageLevel level) { + URI uri = resourceId.toPath().toUri(); + ImmutableList.Builder segments = ImmutableList.builder(); + if (uri.getAuthority() != null && !uri.getAuthority().isEmpty()) { + segments.add(uri.getAuthority()); + } + if (level != LineageLevel.TOP_LEVEL + && uri.getPath() != null + && !uri.getPath().isEmpty() + && !uri.getPath().equals("/")) { + segments.add(uri.getPath()); + } + lineage.add(scheme, segments.build(), "/"); + } + /** An adapter around {@link FSDataInputStream} that implements {@link SeekableByteChannel}. */ private static class HadoopSeekableByteChannel implements SeekableByteChannel { private final FileStatus fileStatus; diff --git a/sdks/java/io/hadoop-file-system/src/test/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java b/sdks/java/io/hadoop-file-system/src/test/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java index c5918a8c9aa6..3642451c002c 100644 --- a/sdks/java/io/hadoop-file-system/src/test/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java +++ b/sdks/java/io/hadoop-file-system/src/test/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java @@ -24,6 +24,9 @@ import static org.hamcrest.Matchers.hasSize; import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import java.io.FileNotFoundException; import java.io.InputStream; @@ -43,6 +46,7 @@ import org.apache.beam.sdk.io.fs.MatchResult; import org.apache.beam.sdk.io.fs.MatchResult.Metadata; import org.apache.beam.sdk.io.fs.MatchResult.Status; +import org.apache.beam.sdk.metrics.Lineage; import org.apache.beam.sdk.testing.ExpectedLogs; import org.apache.beam.sdk.testing.PAssert; import org.apache.beam.sdk.testing.TestPipeline; @@ -481,6 +485,21 @@ public void testReadPipeline() throws Exception { p.run(); } + @Test + public void testReportLineage() { + verifyLineage( + "hdfs://namenode/path/to/file.txt", ImmutableList.of("namenode", "/path/to/file.txt")); + verifyLineage("hdfs://namenode/", ImmutableList.of("namenode")); + verifyLineage("hdfs://namenode", ImmutableList.of("namenode")); + } + + private void verifyLineage(String uri, List expected) { + HadoopResourceId resourceId = new HadoopResourceId(URI.create(uri)); + Lineage mockLineage = mock(Lineage.class); + fileSystem.reportLineage(resourceId, mockLineage); + verify(mockLineage, times(1)).add("hdfs", expected, "/"); + } + private void create(String relativePath, byte[] contents) throws Exception { try (WritableByteChannel channel = fileSystem.create( From 892a0058be0b95ce505f989e94c2cb3f5274916c Mon Sep 17 00:00:00 2001 From: Andrew Kabas Date: Thu, 23 Apr 2026 10:35:36 -0400 Subject: [PATCH 7/8] address review comments --- .../src/main/java/org/apache/beam/sdk/io/FileSystems.java | 3 +++ .../java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java | 5 ----- .../org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java | 3 ++- 3 files changed, 5 insertions(+), 6 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java index e93f4efa8e34..6ae359513862 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java @@ -408,7 +408,10 @@ public ResourceId apply(@Nonnull Metadata input) { *

- Number of directory<=100, report directory names (one level up); * *

- Otherwise, report top level only. + * + *

For internal use only by Beam-provided file-based connectors; not a stable public API. */ + @Internal public static void reportSourceLineage(List resourceIds) { final int maxLineageTargets = 100; if (resourceIds.size() <= maxLineageTargets) { diff --git a/sdks/java/io/hadoop-file-system/src/main/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java b/sdks/java/io/hadoop-file-system/src/main/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java index e0a7793a243e..00a77aa84476 100644 --- a/sdks/java/io/hadoop-file-system/src/main/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java +++ b/sdks/java/io/hadoop-file-system/src/main/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java @@ -337,11 +337,6 @@ protected String getScheme() { return scheme; } - @Override - protected void reportLineage(HadoopResourceId resourceId, Lineage lineage) { - reportLineage(resourceId, lineage, LineageLevel.FILE); - } - @Override protected void reportLineage(HadoopResourceId resourceId, Lineage lineage, LineageLevel level) { URI uri = resourceId.toPath().toUri(); diff --git a/sdks/java/io/hadoop-file-system/src/test/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java b/sdks/java/io/hadoop-file-system/src/test/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java index 3642451c002c..e0b299401c76 100644 --- a/sdks/java/io/hadoop-file-system/src/test/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java +++ b/sdks/java/io/hadoop-file-system/src/test/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java @@ -40,6 +40,7 @@ import java.util.Collections; import java.util.List; import java.util.Objects; +import org.apache.beam.sdk.io.FileSystem; import org.apache.beam.sdk.io.FileSystems; import org.apache.beam.sdk.io.TextIO; import org.apache.beam.sdk.io.fs.CreateOptions.StandardCreateOptions; @@ -496,7 +497,7 @@ public void testReportLineage() { private void verifyLineage(String uri, List expected) { HadoopResourceId resourceId = new HadoopResourceId(URI.create(uri)); Lineage mockLineage = mock(Lineage.class); - fileSystem.reportLineage(resourceId, mockLineage); + fileSystem.reportLineage(resourceId, mockLineage, FileSystem.LineageLevel.FILE); verify(mockLineage, times(1)).add("hdfs", expected, "/"); } From 26c846a45e83c38bd6f9e1f24eb89c2ce10fe728 Mon Sep 17 00:00:00 2001 From: Andrew Kabas Date: Wed, 20 May 2026 10:06:27 -0400 Subject: [PATCH 8/8] Address review comments --- .../apache/beam/sdk/io/FileBasedSource.java | 34 ++++++- .../org/apache/beam/sdk/io/FileSystems.java | 36 ------- .../apache/beam/sdk/io/FileSystemsTest.java | 98 ------------------- .../sdk/io/hadoop/format/HadoopFormatIO.java | 32 +++--- 4 files changed, 50 insertions(+), 150 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSource.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSource.java index 4d65fbbd5e93..8d6e52c64a52 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSource.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSource.java @@ -26,11 +26,12 @@ import java.nio.channels.ReadableByteChannel; import java.nio.channels.SeekableByteChannel; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; import java.util.ListIterator; import java.util.NoSuchElementException; import java.util.concurrent.atomic.AtomicReference; -import java.util.stream.Collectors; +import org.apache.beam.sdk.io.FileSystem.LineageLevel; import org.apache.beam.sdk.io.fs.EmptyMatchTreatment; import org.apache.beam.sdk.io.fs.MatchResult; import org.apache.beam.sdk.io.fs.MatchResult.Metadata; @@ -317,10 +318,35 @@ public final List> split( } } + /** + * Report source Lineage. Due to the size limit of Beam metrics, report full file name or only dir + * depend on the number of files. + * + *

- Number of files<=100, report full file paths; + * + *

- Number of directory<=100, report directory names (one level up); + * + *

- Otherwise, report top level only. + */ private static void reportSourceLineage(List expandedFiles) { - List resourceIds = - expandedFiles.stream().map(Metadata::resourceId).collect(Collectors.toList()); - FileSystems.reportSourceLineage(resourceIds); + if (expandedFiles.size() <= 100) { + for (Metadata metadata : expandedFiles) { + FileSystems.reportSourceLineage(metadata.resourceId()); + } + } else { + HashSet uniqueDirs = new HashSet<>(); + for (Metadata metadata : expandedFiles) { + ResourceId dir = metadata.resourceId().getCurrentDirectory(); + uniqueDirs.add(dir); + if (uniqueDirs.size() > 100) { + FileSystems.reportSourceLineage(dir, LineageLevel.TOP_LEVEL); + return; + } + } + for (ResourceId uniqueDir : uniqueDirs) { + FileSystems.reportSourceLineage(uniqueDir); + } + } } /** diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java index 6ae359513862..6133ca9fdb39 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java @@ -29,7 +29,6 @@ import java.util.ArrayList; import java.util.Collection; import java.util.Collections; -import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Map.Entry; @@ -399,41 +398,6 @@ public ResourceId apply(@Nonnull Metadata input) { .delete(resourceIdsToDelete); } - /** - * Report source {@link Lineage} metrics for multiple resource ids. Due to the size limit of Beam - * metrics, report full file name or only dir depend on the number of files. - * - *

- Number of files<=100, report full file paths; - * - *

- Number of directory<=100, report directory names (one level up); - * - *

- Otherwise, report top level only. - * - *

For internal use only by Beam-provided file-based connectors; not a stable public API. - */ - @Internal - public static void reportSourceLineage(List resourceIds) { - final int maxLineageTargets = 100; - if (resourceIds.size() <= maxLineageTargets) { - for (ResourceId resourceId : resourceIds) { - FileSystems.reportSourceLineage(resourceId); - } - } else { - HashSet uniqueDirs = new HashSet<>(); - for (ResourceId resourceId : resourceIds) { - ResourceId dir = resourceId.getCurrentDirectory(); - uniqueDirs.add(dir); - if (uniqueDirs.size() > maxLineageTargets) { - FileSystems.reportSourceLineage(dir, LineageLevel.TOP_LEVEL); - return; - } - } - for (ResourceId uniqueDir : uniqueDirs) { - FileSystems.reportSourceLineage(uniqueDir); - } - } - } - /** Report source {@link Lineage} metrics for resource id. */ public static void reportSourceLineage(ResourceId resourceId) { reportSourceLineage(resourceId, LineageLevel.FILE); diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileSystemsTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileSystemsTest.java index bd7cb413739b..34567309c7d0 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileSystemsTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileSystemsTest.java @@ -23,11 +23,8 @@ import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; import static org.junit.Assume.assumeFalse; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -38,9 +35,7 @@ import java.nio.file.NoSuchFileException; import java.nio.file.Path; import java.nio.file.Paths; -import java.util.ArrayList; import java.util.List; -import org.apache.beam.sdk.io.FileSystem.LineageLevel; import org.apache.beam.sdk.io.fs.CreateOptions; import org.apache.beam.sdk.io.fs.MatchResult; import org.apache.beam.sdk.io.fs.MoveOptions; @@ -59,8 +54,6 @@ import org.junit.rules.TemporaryFolder; import org.junit.runner.RunWith; import org.junit.runners.JUnit4; -import org.mockito.MockedStatic; -import org.mockito.Mockito; /** Tests for {@link FileSystems}. */ @RunWith(JUnit4.class) @@ -344,97 +337,6 @@ public void testMatchNewDirectory() { } } - @Test - public void testReportSourceLineageFewFiles() { - List resources = new ArrayList<>(); - for (int i = 0; i < 5; i++) { - resources.add(mockResourceId("/dir/file" + i, "/dir/")); - } - - try (MockedStatic mocked = - Mockito.mockStatic(FileSystems.class, Mockito.CALLS_REAL_METHODS)) { - mocked - .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class))) - .thenAnswer(inv -> null); - mocked - .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class), any())) - .thenAnswer(inv -> null); - - FileSystems.reportSourceLineage(resources); - - for (ResourceId r : resources) { - mocked.verify(() -> FileSystems.reportSourceLineage(r)); - } - } - } - - @Test - public void testReportSourceLineageManyFilesFewDirs() { - ResourceId dir = mockResourceId("/dir/", null); - List resources = new ArrayList<>(); - for (int i = 0; i < 150; i++) { - resources.add(mockResourceId("/dir/file" + i, dir)); - } - - try (MockedStatic mocked = - Mockito.mockStatic(FileSystems.class, Mockito.CALLS_REAL_METHODS)) { - mocked - .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class))) - .thenAnswer(inv -> null); - mocked - .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class), any())) - .thenAnswer(inv -> null); - - FileSystems.reportSourceLineage(resources); - - // Should report the unique directory, not individual files - mocked.verify(() -> FileSystems.reportSourceLineage(dir)); - mocked.verify( - () -> FileSystems.reportSourceLineage(any(ResourceId.class), any(LineageLevel.class)), - never()); - } - } - - @Test - public void testReportSourceLineageManyFilesManyDirs() { - List resources = new ArrayList<>(); - for (int i = 0; i < 150; i++) { - ResourceId dir = mockResourceId("/dir" + i + "/", null); - resources.add(mockResourceId("/dir" + i + "/file", dir)); - } - - try (MockedStatic mocked = - Mockito.mockStatic(FileSystems.class, Mockito.CALLS_REAL_METHODS)) { - mocked - .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class))) - .thenAnswer(inv -> null); - mocked - .when(() -> FileSystems.reportSourceLineage(any(ResourceId.class), any())) - .thenAnswer(inv -> null); - - FileSystems.reportSourceLineage(resources); - - // Should fall back to TOP_LEVEL reporting - mocked.verify( - () -> FileSystems.reportSourceLineage(any(ResourceId.class), eq(LineageLevel.TOP_LEVEL))); - // Should not report individual files or directories at FILE level - mocked.verify(() -> FileSystems.reportSourceLineage(any(ResourceId.class)), never()); - } - } - - private ResourceId mockResourceId(String path, Object dir) { - ResourceId resourceId = mock(ResourceId.class); - when(resourceId.toString()).thenReturn(path); - if (dir instanceof String) { - ResourceId dirId = mock(ResourceId.class); - when(dirId.toString()).thenReturn((String) dir); - when(resourceId.getCurrentDirectory()).thenReturn(dirId); - } else if (dir instanceof ResourceId) { - when(resourceId.getCurrentDirectory()).thenReturn((ResourceId) dir); - } - return resourceId; - } - private static List toResourceIds(List paths, final boolean isDirectory) { return FluentIterable.from(paths) .transform(path -> (ResourceId) LocalResourceId.fromPath(path, isDirectory)) diff --git a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java index 58780a9eb66c..501040928e41 100644 --- a/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java +++ b/sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java @@ -34,6 +34,7 @@ import java.lang.reflect.InvocationTargetException; import java.math.BigDecimal; import java.math.BigInteger; +import java.net.URI; import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; @@ -52,10 +53,9 @@ import org.apache.beam.sdk.coders.CoderRegistry; import org.apache.beam.sdk.coders.KvCoder; import org.apache.beam.sdk.io.BoundedSource; -import org.apache.beam.sdk.io.FileSystems; -import org.apache.beam.sdk.io.fs.ResourceId; import org.apache.beam.sdk.io.hadoop.SerializableConfiguration; import org.apache.beam.sdk.io.hadoop.WritableCoder; +import org.apache.beam.sdk.metrics.Lineage; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.transforms.Combine; import org.apache.beam.sdk.transforms.Create; @@ -748,17 +748,25 @@ public List>> split(long desiredBundleSizeBytes, Pipeline .collect(Collectors.toList()); } - /** Report only file-based sources. */ private void reportSourceLineage(final List inputSplits) { - List fileResources = - inputSplits.stream() - .map(SerializableSplit::getSplit) - .filter(FileSplit.class::isInstance) - .map(FileSplit.class::cast) - .map(fileSplit -> FileSystems.matchNewResource(fileSplit.getPath().toString(), false)) - .collect(Collectors.toList()); - - FileSystems.reportSourceLineage(fileResources); + for (SerializableSplit serializableSplit : inputSplits) { + InputSplit split = serializableSplit.getSplit(); + if (split instanceof FileSplit) { + URI uri = ((FileSplit) split).getPath().toUri(); + String scheme = uri.getScheme(); + if (scheme == null) { + continue; + } + ImmutableList.Builder segments = ImmutableList.builder(); + if (uri.getAuthority() != null && !uri.getAuthority().isEmpty()) { + segments.add(uri.getAuthority()); + } + if (uri.getPath() != null && !uri.getPath().isEmpty() && !uri.getPath().equals("/")) { + segments.add(uri.getPath()); + } + Lineage.getSources().add(scheme, segments.build(), "/"); + } + } } @Override