Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import org.apache.beam.sdk.lineage.LineageOptions;
import org.apache.beam.sdk.metrics.Metrics.MetricsFlag;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
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.base.Splitter;
import org.checkerframework.checker.nullness.qual.Nullable;
Expand Down Expand Up @@ -122,16 +123,24 @@ private static Lineage createLineage(PipelineOptions options, LineageDirection d

/** {@link Lineage} representing sources and optionally side inputs. */
public static Lineage getSources() {
return checkNotNull(
sources,
"Lineage not initialized. FileSystems.setDefaultPipelineOptions must be called first.");
Lineage localSources = sources;
if (localSources == null) {
return createDefaultLineage(LineageDirection.SOURCE);
}
return localSources;
}

/** {@link Lineage} representing sinks. */
public static Lineage getSinks() {
return checkNotNull(
sinks,
"Lineage not initialized. FileSystems.setDefaultPipelineOptions must be called first.");
Lineage localSinks = sinks;
if (localSinks == null) {
return createDefaultLineage(LineageDirection.SINK);
}
return localSinks;
}

private static Lineage createDefaultLineage(LineageDirection direction) {
return createLineage(PipelineOptionsFactory.create(), direction);
}

@VisibleForTesting
Expand Down
Loading