Skip to content

Commit 26c846a

Browse files
committed
Address review comments
1 parent 892a005 commit 26c846a

4 files changed

Lines changed: 50 additions & 150 deletions

File tree

sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileBasedSource.java

Lines changed: 30 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -26,11 +26,12 @@
2626
import java.nio.channels.ReadableByteChannel;
2727
import java.nio.channels.SeekableByteChannel;
2828
import java.util.ArrayList;
29+
import java.util.HashSet;
2930
import java.util.List;
3031
import java.util.ListIterator;
3132
import java.util.NoSuchElementException;
3233
import java.util.concurrent.atomic.AtomicReference;
33-
import java.util.stream.Collectors;
34+
import org.apache.beam.sdk.io.FileSystem.LineageLevel;
3435
import org.apache.beam.sdk.io.fs.EmptyMatchTreatment;
3536
import org.apache.beam.sdk.io.fs.MatchResult;
3637
import org.apache.beam.sdk.io.fs.MatchResult.Metadata;
@@ -317,10 +318,35 @@ public final List<? extends FileBasedSource<T>> split(
317318
}
318319
}
319320

321+
/**
322+
* Report source Lineage. Due to the size limit of Beam metrics, report full file name or only dir
323+
* depend on the number of files.
324+
*
325+
* <p>- Number of files<=100, report full file paths;
326+
*
327+
* <p>- Number of directory<=100, report directory names (one level up);
328+
*
329+
* <p>- Otherwise, report top level only.
330+
*/
320331
private static void reportSourceLineage(List<Metadata> expandedFiles) {
321-
List<ResourceId> resourceIds =
322-
expandedFiles.stream().map(Metadata::resourceId).collect(Collectors.toList());
323-
FileSystems.reportSourceLineage(resourceIds);
332+
if (expandedFiles.size() <= 100) {
333+
for (Metadata metadata : expandedFiles) {
334+
FileSystems.reportSourceLineage(metadata.resourceId());
335+
}
336+
} else {
337+
HashSet<ResourceId> uniqueDirs = new HashSet<>();
338+
for (Metadata metadata : expandedFiles) {
339+
ResourceId dir = metadata.resourceId().getCurrentDirectory();
340+
uniqueDirs.add(dir);
341+
if (uniqueDirs.size() > 100) {
342+
FileSystems.reportSourceLineage(dir, LineageLevel.TOP_LEVEL);
343+
return;
344+
}
345+
}
346+
for (ResourceId uniqueDir : uniqueDirs) {
347+
FileSystems.reportSourceLineage(uniqueDir);
348+
}
349+
}
324350
}
325351

326352
/**

sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileSystems.java

Lines changed: 0 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,6 @@
2929
import java.util.ArrayList;
3030
import java.util.Collection;
3131
import java.util.Collections;
32-
import java.util.HashSet;
3332
import java.util.List;
3433
import java.util.Map;
3534
import java.util.Map.Entry;
@@ -399,41 +398,6 @@ public ResourceId apply(@Nonnull Metadata input) {
399398
.delete(resourceIdsToDelete);
400399
}
401400

402-
/**
403-
* Report source {@link Lineage} metrics for multiple resource ids. Due to the size limit of Beam
404-
* metrics, report full file name or only dir depend on the number of files.
405-
*
406-
* <p>- Number of files<=100, report full file paths;
407-
*
408-
* <p>- Number of directory<=100, report directory names (one level up);
409-
*
410-
* <p>- Otherwise, report top level only.
411-
*
412-
* <p>For internal use only by Beam-provided file-based connectors; not a stable public API.
413-
*/
414-
@Internal
415-
public static void reportSourceLineage(List<ResourceId> resourceIds) {
416-
final int maxLineageTargets = 100;
417-
if (resourceIds.size() <= maxLineageTargets) {
418-
for (ResourceId resourceId : resourceIds) {
419-
FileSystems.reportSourceLineage(resourceId);
420-
}
421-
} else {
422-
HashSet<ResourceId> uniqueDirs = new HashSet<>();
423-
for (ResourceId resourceId : resourceIds) {
424-
ResourceId dir = resourceId.getCurrentDirectory();
425-
uniqueDirs.add(dir);
426-
if (uniqueDirs.size() > maxLineageTargets) {
427-
FileSystems.reportSourceLineage(dir, LineageLevel.TOP_LEVEL);
428-
return;
429-
}
430-
}
431-
for (ResourceId uniqueDir : uniqueDirs) {
432-
FileSystems.reportSourceLineage(uniqueDir);
433-
}
434-
}
435-
}
436-
437401
/** Report source {@link Lineage} metrics for resource id. */
438402
public static void reportSourceLineage(ResourceId resourceId) {
439403
reportSourceLineage(resourceId, LineageLevel.FILE);

sdks/java/core/src/test/java/org/apache/beam/sdk/io/FileSystemsTest.java

Lines changed: 0 additions & 98 deletions
Original file line numberDiff line numberDiff line change
@@ -23,11 +23,8 @@
2323
import static org.junit.Assert.assertFalse;
2424
import static org.junit.Assert.assertTrue;
2525
import static org.junit.Assume.assumeFalse;
26-
import static org.mockito.ArgumentMatchers.any;
27-
import static org.mockito.ArgumentMatchers.eq;
2826
import static org.mockito.Mockito.doThrow;
2927
import static org.mockito.Mockito.mock;
30-
import static org.mockito.Mockito.never;
3128
import static org.mockito.Mockito.verify;
3229
import static org.mockito.Mockito.when;
3330

@@ -38,9 +35,7 @@
3835
import java.nio.file.NoSuchFileException;
3936
import java.nio.file.Path;
4037
import java.nio.file.Paths;
41-
import java.util.ArrayList;
4238
import java.util.List;
43-
import org.apache.beam.sdk.io.FileSystem.LineageLevel;
4439
import org.apache.beam.sdk.io.fs.CreateOptions;
4540
import org.apache.beam.sdk.io.fs.MatchResult;
4641
import org.apache.beam.sdk.io.fs.MoveOptions;
@@ -59,8 +54,6 @@
5954
import org.junit.rules.TemporaryFolder;
6055
import org.junit.runner.RunWith;
6156
import org.junit.runners.JUnit4;
62-
import org.mockito.MockedStatic;
63-
import org.mockito.Mockito;
6457

6558
/** Tests for {@link FileSystems}. */
6659
@RunWith(JUnit4.class)
@@ -344,97 +337,6 @@ public void testMatchNewDirectory() {
344337
}
345338
}
346339

347-
@Test
348-
public void testReportSourceLineageFewFiles() {
349-
List<ResourceId> resources = new ArrayList<>();
350-
for (int i = 0; i < 5; i++) {
351-
resources.add(mockResourceId("/dir/file" + i, "/dir/"));
352-
}
353-
354-
try (MockedStatic<FileSystems> mocked =
355-
Mockito.mockStatic(FileSystems.class, Mockito.CALLS_REAL_METHODS)) {
356-
mocked
357-
.when(() -> FileSystems.reportSourceLineage(any(ResourceId.class)))
358-
.thenAnswer(inv -> null);
359-
mocked
360-
.when(() -> FileSystems.reportSourceLineage(any(ResourceId.class), any()))
361-
.thenAnswer(inv -> null);
362-
363-
FileSystems.reportSourceLineage(resources);
364-
365-
for (ResourceId r : resources) {
366-
mocked.verify(() -> FileSystems.reportSourceLineage(r));
367-
}
368-
}
369-
}
370-
371-
@Test
372-
public void testReportSourceLineageManyFilesFewDirs() {
373-
ResourceId dir = mockResourceId("/dir/", null);
374-
List<ResourceId> resources = new ArrayList<>();
375-
for (int i = 0; i < 150; i++) {
376-
resources.add(mockResourceId("/dir/file" + i, dir));
377-
}
378-
379-
try (MockedStatic<FileSystems> mocked =
380-
Mockito.mockStatic(FileSystems.class, Mockito.CALLS_REAL_METHODS)) {
381-
mocked
382-
.when(() -> FileSystems.reportSourceLineage(any(ResourceId.class)))
383-
.thenAnswer(inv -> null);
384-
mocked
385-
.when(() -> FileSystems.reportSourceLineage(any(ResourceId.class), any()))
386-
.thenAnswer(inv -> null);
387-
388-
FileSystems.reportSourceLineage(resources);
389-
390-
// Should report the unique directory, not individual files
391-
mocked.verify(() -> FileSystems.reportSourceLineage(dir));
392-
mocked.verify(
393-
() -> FileSystems.reportSourceLineage(any(ResourceId.class), any(LineageLevel.class)),
394-
never());
395-
}
396-
}
397-
398-
@Test
399-
public void testReportSourceLineageManyFilesManyDirs() {
400-
List<ResourceId> resources = new ArrayList<>();
401-
for (int i = 0; i < 150; i++) {
402-
ResourceId dir = mockResourceId("/dir" + i + "/", null);
403-
resources.add(mockResourceId("/dir" + i + "/file", dir));
404-
}
405-
406-
try (MockedStatic<FileSystems> mocked =
407-
Mockito.mockStatic(FileSystems.class, Mockito.CALLS_REAL_METHODS)) {
408-
mocked
409-
.when(() -> FileSystems.reportSourceLineage(any(ResourceId.class)))
410-
.thenAnswer(inv -> null);
411-
mocked
412-
.when(() -> FileSystems.reportSourceLineage(any(ResourceId.class), any()))
413-
.thenAnswer(inv -> null);
414-
415-
FileSystems.reportSourceLineage(resources);
416-
417-
// Should fall back to TOP_LEVEL reporting
418-
mocked.verify(
419-
() -> FileSystems.reportSourceLineage(any(ResourceId.class), eq(LineageLevel.TOP_LEVEL)));
420-
// Should not report individual files or directories at FILE level
421-
mocked.verify(() -> FileSystems.reportSourceLineage(any(ResourceId.class)), never());
422-
}
423-
}
424-
425-
private ResourceId mockResourceId(String path, Object dir) {
426-
ResourceId resourceId = mock(ResourceId.class);
427-
when(resourceId.toString()).thenReturn(path);
428-
if (dir instanceof String) {
429-
ResourceId dirId = mock(ResourceId.class);
430-
when(dirId.toString()).thenReturn((String) dir);
431-
when(resourceId.getCurrentDirectory()).thenReturn(dirId);
432-
} else if (dir instanceof ResourceId) {
433-
when(resourceId.getCurrentDirectory()).thenReturn((ResourceId) dir);
434-
}
435-
return resourceId;
436-
}
437-
438340
private static List<ResourceId> toResourceIds(List<Path> paths, final boolean isDirectory) {
439341
return FluentIterable.from(paths)
440342
.transform(path -> (ResourceId) LocalResourceId.fromPath(path, isDirectory))

sdks/java/io/hadoop-format/src/main/java/org/apache/beam/sdk/io/hadoop/format/HadoopFormatIO.java

Lines changed: 20 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
import java.lang.reflect.InvocationTargetException;
3535
import java.math.BigDecimal;
3636
import java.math.BigInteger;
37+
import java.net.URI;
3738
import java.util.ArrayList;
3839
import java.util.Arrays;
3940
import java.util.HashMap;
@@ -52,10 +53,9 @@
5253
import org.apache.beam.sdk.coders.CoderRegistry;
5354
import org.apache.beam.sdk.coders.KvCoder;
5455
import org.apache.beam.sdk.io.BoundedSource;
55-
import org.apache.beam.sdk.io.FileSystems;
56-
import org.apache.beam.sdk.io.fs.ResourceId;
5756
import org.apache.beam.sdk.io.hadoop.SerializableConfiguration;
5857
import org.apache.beam.sdk.io.hadoop.WritableCoder;
58+
import org.apache.beam.sdk.metrics.Lineage;
5959
import org.apache.beam.sdk.options.PipelineOptions;
6060
import org.apache.beam.sdk.transforms.Combine;
6161
import org.apache.beam.sdk.transforms.Create;
@@ -748,17 +748,25 @@ public List<BoundedSource<KV<K, V>>> split(long desiredBundleSizeBytes, Pipeline
748748
.collect(Collectors.toList());
749749
}
750750

751-
/** Report only file-based sources. */
752751
private void reportSourceLineage(final List<SerializableSplit> inputSplits) {
753-
List<ResourceId> fileResources =
754-
inputSplits.stream()
755-
.map(SerializableSplit::getSplit)
756-
.filter(FileSplit.class::isInstance)
757-
.map(FileSplit.class::cast)
758-
.map(fileSplit -> FileSystems.matchNewResource(fileSplit.getPath().toString(), false))
759-
.collect(Collectors.toList());
760-
761-
FileSystems.reportSourceLineage(fileResources);
752+
for (SerializableSplit serializableSplit : inputSplits) {
753+
InputSplit split = serializableSplit.getSplit();
754+
if (split instanceof FileSplit) {
755+
URI uri = ((FileSplit) split).getPath().toUri();
756+
String scheme = uri.getScheme();
757+
if (scheme == null) {
758+
continue;
759+
}
760+
ImmutableList.Builder<String> segments = ImmutableList.builder();
761+
if (uri.getAuthority() != null && !uri.getAuthority().isEmpty()) {
762+
segments.add(uri.getAuthority());
763+
}
764+
if (uri.getPath() != null && !uri.getPath().isEmpty() && !uri.getPath().equals("/")) {
765+
segments.add(uri.getPath());
766+
}
767+
Lineage.getSources().add(scheme, segments.build(), "/");
768+
}
769+
}
762770
}
763771

764772
@Override

0 commit comments

Comments
 (0)