Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:Datahub project Datahub RddPathUtils

From Leeroopedia
Revision as of 14:43, 16 February 2026 by Admin (talk | contribs) (Auto-imported from implementations/Datahub_project_Datahub_RddPathUtils.md)
(diff) ← Older revision | Latest revision (diff) | Newer revision → (diff)


Knowledge Sources
Domains Spark_Lineage, OpenLineage
Last Updated 2026-02-10 00:00 GMT

Overview

RddPathUtils is a utility class in the io.openlineage.spark.agent.util package that extracts filesystem paths from Spark RDD (Resilient Distributed Dataset) nodes. It provides the low-level mechanism for discovering input dataset locations from the RDD lineage graph, which is essential for capturing lineage from RDD-based Spark jobs and from the underlying RDDs of DataFrame/Dataset operations.

The class uses a strategy pattern with an internal RddPathExtractor interface. Four concrete extractors handle specific RDD types, and an UnknownRDDExtractor serves as a fallback. The extraction is recursive for MapPartitionsRDD, traversing up the RDD lineage chain until a concrete path-bearing RDD is found.

Source file: metadata-integration/java/acryl-spark-lineage/src/main/java/io/openlineage/spark/agent/util/RddPathUtils.java (192 lines)

Code Reference

Class Declaration

@Slf4j
public class RddPathUtils {

Main Entry Point

findRDDPaths

public static Stream<Path> findRDDPaths(RDD rdd)

Accepts any RDD and returns a stream of Hadoop Path objects representing the input data locations. Iterates through the registered extractors in order, selects the first one whose isDefinedAt returns true, and invokes its extract method. Falls back to UnknownRDDExtractor if no extractor matches.

RddPathExtractor Interface

interface RddPathExtractor<T extends RDD> {
    boolean isDefinedAt(Object rdd);
    Stream<Path> extract(T rdd);
}

Concrete Extractors

HadoopRDDExtractor

static class HadoopRDDExtractor implements RddPathExtractor<HadoopRDD>

Handles HadoopRDD instances. Retrieves input paths from FileInputFormat.getInputPaths using the RDD's job configuration, then resolves each path to its directory via PlanUtils.getDirectoryPath.

FileScanRDDExtractor

static class FileScanRDDExtractor implements RddPathExtractor<FileScanRDD>

Handles FileScanRDD instances. Iterates over file partitions and their files. Accounts for the Spark version difference:

  • Spark 3.4+: filePath() returns SparkPath; calls toPath() via reflection.
  • Spark < 3.4: filePath() returns String; parses it directly.

Returns the parent directory of each file path.

MapPartitionsRDDExtractor

static class MapPartitionsRDDExtractor implements RddPathExtractor<MapPartitionsRDD>

Handles MapPartitionsRDD by recursively calling findRDDPaths on the parent RDD (rdd.prev()), traversing the RDD lineage chain.

ParallelCollectionRDDExtractor

static class ParallelCollectionRDDExtractor implements RddPathExtractor<ParallelCollectionRDD>

Handles ParallelCollectionRDD (typically from sc.parallelize). Uses reflection to read the data field and extracts paths from:

  • Seq<Tuple2> - extracts paths from the first element of each tuple (limited to 1000 entries)
  • ArrayBuffer - converts elements to paths directly

UnknownRDDExtractor

static class UnknownRDDExtractor implements RddPathExtractor<RDD>

Fallback extractor that matches any RDD type. Returns an empty stream and logs at debug level.

Helper Method

private static Path parentOf(String path)

Safely parses a string into a Hadoop Path and returns its parent directory, returning null on failure.

I/O Contract

Direction Type Description
Input RDD Any Spark RDD instance from which to extract filesystem paths.
Output Stream<Path> A stream of Hadoop Path objects representing the input data directories found in the RDD lineage.

Usage Examples

Extracting paths from an RDD in a Spark job:

RDD<?> rdd = sparkContext.hadoopFile("/data/input", ...);
Stream<Path> paths = RddPathUtils.findRDDPaths(rdd);
paths.forEach(path -> {
    DatasetIdentifier id = PathUtils.fromPath(path);
    // Use id for lineage tracking
});

For a MapPartitionsRDD (e.g., after a map() transformation), the extractor recursively traverses to the source:

// rdd = hadoopRDD.map(transformFunction)
// MapPartitionsRDDExtractor -> calls findRDDPaths(hadoopRDD)
// HadoopRDDExtractor -> returns input paths from FileInputFormat

Related Pages

Page Connections

Double-click a node to navigate. Hold to expand connections.
Principle
Implementation
Heuristic
Environment