Skip to content

Commit fc948dd

Browse files
committed
Implement lineage in HDFS and add tests
1 parent 3fe632c commit fc948dd

3 files changed

Lines changed: 139 additions & 0 deletions

File tree

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

Lines changed: 98 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,8 +23,11 @@
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;
2628
import static org.mockito.Mockito.doThrow;
2729
import static org.mockito.Mockito.mock;
30+
import static org.mockito.Mockito.never;
2831
import static org.mockito.Mockito.verify;
2932
import static org.mockito.Mockito.when;
3033

@@ -35,7 +38,9 @@
3538
import java.nio.file.NoSuchFileException;
3639
import java.nio.file.Path;
3740
import java.nio.file.Paths;
41+
import java.util.ArrayList;
3842
import java.util.List;
43+
import org.apache.beam.sdk.io.FileSystem.LineageLevel;
3944
import org.apache.beam.sdk.io.fs.CreateOptions;
4045
import org.apache.beam.sdk.io.fs.MatchResult;
4146
import org.apache.beam.sdk.io.fs.MoveOptions;
@@ -54,6 +59,8 @@
5459
import org.junit.rules.TemporaryFolder;
5560
import org.junit.runner.RunWith;
5661
import org.junit.runners.JUnit4;
62+
import org.mockito.MockedStatic;
63+
import org.mockito.Mockito;
5764

5865
/** Tests for {@link FileSystems}. */
5966
@RunWith(JUnit4.class)
@@ -337,6 +344,97 @@ public void testMatchNewDirectory() {
337344
}
338345
}
339346

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+
340438
private static List<ResourceId> toResourceIds(List<Path> paths, final boolean isDirectory) {
341439
return FluentIterable.from(paths)
342440
.transform(path -> (ResourceId) LocalResourceId.fromPath(path, isDirectory))

sdks/java/io/hadoop-file-system/src/main/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystem.java

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838
import org.apache.beam.sdk.io.fs.MatchResult.Metadata;
3939
import org.apache.beam.sdk.io.fs.MatchResult.Status;
4040
import org.apache.beam.sdk.io.fs.MoveOptions;
41+
import org.apache.beam.sdk.metrics.Lineage;
4142
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
4243
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
4344
import org.apache.hadoop.conf.Configuration;
@@ -336,6 +337,27 @@ protected String getScheme() {
336337
return scheme;
337338
}
338339

340+
@Override
341+
protected void reportLineage(HadoopResourceId resourceId, Lineage lineage) {
342+
reportLineage(resourceId, lineage, LineageLevel.FILE);
343+
}
344+
345+
@Override
346+
protected void reportLineage(HadoopResourceId resourceId, Lineage lineage, LineageLevel level) {
347+
URI uri = resourceId.toPath().toUri();
348+
ImmutableList.Builder<String> segments = ImmutableList.builder();
349+
if (uri.getAuthority() != null && !uri.getAuthority().isEmpty()) {
350+
segments.add(uri.getAuthority());
351+
}
352+
if (level != LineageLevel.TOP_LEVEL
353+
&& uri.getPath() != null
354+
&& !uri.getPath().isEmpty()
355+
&& !uri.getPath().equals("/")) {
356+
segments.add(uri.getPath());
357+
}
358+
lineage.add(scheme, segments.build(), "/");
359+
}
360+
339361
/** An adapter around {@link FSDataInputStream} that implements {@link SeekableByteChannel}. */
340362
private static class HadoopSeekableByteChannel implements SeekableByteChannel {
341363
private final FileStatus fileStatus;

sdks/java/io/hadoop-file-system/src/test/java/org/apache/beam/sdk/io/hdfs/HadoopFileSystemTest.java

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,9 @@
2424
import static org.hamcrest.Matchers.hasSize;
2525
import static org.junit.Assert.assertArrayEquals;
2626
import static org.junit.Assert.assertEquals;
27+
import static org.mockito.Mockito.mock;
28+
import static org.mockito.Mockito.times;
29+
import static org.mockito.Mockito.verify;
2730

2831
import java.io.FileNotFoundException;
2932
import java.io.InputStream;
@@ -43,6 +46,7 @@
4346
import org.apache.beam.sdk.io.fs.MatchResult;
4447
import org.apache.beam.sdk.io.fs.MatchResult.Metadata;
4548
import org.apache.beam.sdk.io.fs.MatchResult.Status;
49+
import org.apache.beam.sdk.metrics.Lineage;
4650
import org.apache.beam.sdk.testing.ExpectedLogs;
4751
import org.apache.beam.sdk.testing.PAssert;
4852
import org.apache.beam.sdk.testing.TestPipeline;
@@ -481,6 +485,21 @@ public void testReadPipeline() throws Exception {
481485
p.run();
482486
}
483487

488+
@Test
489+
public void testReportLineage() {
490+
verifyLineage(
491+
"hdfs://namenode/path/to/file.txt", ImmutableList.of("namenode", "/path/to/file.txt"));
492+
verifyLineage("hdfs://namenode/", ImmutableList.of("namenode"));
493+
verifyLineage("hdfs://namenode", ImmutableList.of("namenode"));
494+
}
495+
496+
private void verifyLineage(String uri, List<String> expected) {
497+
HadoopResourceId resourceId = new HadoopResourceId(URI.create(uri));
498+
Lineage mockLineage = mock(Lineage.class);
499+
fileSystem.reportLineage(resourceId, mockLineage);
500+
verify(mockLineage, times(1)).add("hdfs", expected, "/");
501+
}
502+
484503
private void create(String relativePath, byte[] contents) throws Exception {
485504
try (WritableByteChannel channel =
486505
fileSystem.create(

0 commit comments

Comments
 (0)