Column Lineage
SDLB knows which DataObject an Action reads and writes, see DAG. This is lineage on the level of whole datasets. Column level lineage goes one step further and tells from which columns of which input DataObjects a column of an output DataObject is created, and how. It answers questions like
- which transformations have been applied to create a specific column of a reporting table
- where and how is a column of a source system used further downstream
Exporting the column lineage
The column lineage is analyzed in the init-phase of a dry-run, so no data is read and nothing is written:
sdlb --config config/ --feed-sel '.*' --test dry-run-with-lineage-export
The export writes one Json document per output DataObject to global.dataObjectsSchemaSource, the same
location the schema export uses, see Schema:
global {
dataObjectsSchemaSource = "file:./schema"
}
Format
A document contains the Action which created the DataObject and the lineage of its columns, in the format of the OpenLineage column lineage facet. Using an established format makes the export consumable by existing lineage tools, e.g. Marquez or DataHub, without writing a converter first.
{
"actionId": "computeCityStatistics",
"dataObjectId": "cityStatistics",
"columnLineage": {
"fields": {
"city": {
"inputFields": [
{
"namespace": "sdlb",
"name": "cities",
"field": "name",
"transformations": [{"type": "DIRECT", "subtype": "IDENTITY", "masking": false}]
}
]
},
"population": {
"inputFields": [
{
"namespace": "sdlb",
"name": "cities",
"field": "inhabitants",
"transformations": [
{"type": "DIRECT", "subtype": "TRANSFORMATION", "description": "sum(inhabitants)", "masking": false}
]
}
]
}
}
}
}
name is the id of the input DataObject, and field the name of its column. Note that namespace and name
do not identify the physical dataset, as the same DataObject can be backed by different storage locations in
different environments.
subtype is IDENTITY if the value of the input column is taken over unchanged, e.g. by a select or a rename,
and TRANSFORMATION if it is modified. For a transformation, description holds the expression which creates
the column, cut off if it is very long.
Columns without a source, and columns SDLB could not trace
A column which is not created from an input column at all, e.g. a constant or count(*), is exported with an
empty inputFields list and the expression which creates it. A column whose lineage could not be traced
back completely is listed in unresolvedFields instead, and left out of fields:
{
"actionId": "computeCityStatistics",
"dataObjectId": "cityStatistics",
"columnLineage": {
"fields": {
"loadedAt": {"inputFields": [], "expression": "current_timestamp()"}
}
},
"unresolvedFields": ["externalRating"]
}
The distinction matters when reading the export: an empty inputFields means SDLB knows the column has no
source column, while unresolvedFields means SDLB could not find out. A column missing from both never
existed in the output DataObject.
The lineage of one run covers one Action each. End-to-end lineage over a whole pipeline is assembled by following the columns from Action to Action: the input fields of a DataObject written by one Action are the output columns of the DataObject read by the next one.
Debugging unresolved columns
If many columns end up in unresolvedFields, enable the debug output to see why they could not be traced
back. It is switched on with the environment parameter columnLineageDebug, e.g. as a Java system property:
sdlb --config config/ --feed-sel '.*' --test dry-run-with-lineage-export -Dsdl.columnLineageDebug=true
It can also be set as environment variable SDL_COLUMN_LINEAGE_DEBUG=true or in the configuration under
global.environment.columnLineageDebug. For every DataObject with unresolved columns SDLB then writes an
additional text file <dataObjectId>.lineage-debug.txt next to the exported lineage, holding the columns of
the input DataObjects, every column which could not be traced back, and the plan the lineage was read from:
Action selectCities -> DataObject tgt1 (engine Spark)
Inputs:
src1: name#21, country#22
Unresolved columns:
country
country#26 <- _2#24
produced by LocalRelation
LocalRelation [_1#23, _2#24]
Analyzed plan:
Project [name#21, country#26]
+- Join Inner, (country#22 = country#26)
:- SubqueryAlias a
: +- Project [_1#14 AS name#21, _2#15 AS country#22]
: +- LocalRelation [_1#14, _2#15]
+- SubqueryAlias b
+- Project [_1#23 AS name#25, _2#24 AS country#26]
+- LocalRelation [_1#23, _2#24]
How to read it:
- the line below an unresolved column is the path SDLB followed, from the output column down to the column it
could not trace back any further, e.g.
country#26 <- _2#24. produced bynames the type of the plan node which created that column. If that node type is one SDLB does not handle yet, this is the missing piece of the analysis - please report it with the file.not used by the output DataFrameafter an input DataObject means that none of the columns of the output DataFrame is derived from these input columns. Either the transformation really does not use them, or it re-created them, e.g. by reading the data itself, which breaks the lineage.
A DataObject whose columns are all traced back gets no debug file, so the files which are there are the ones to look at. The debug output describes engine internals, its format is not stable, and it is meant to be read by a developer. Leave it switched off for regular exports.
Caching
An Action can hand its output DataFrame to the next Action instead of letting it read the DataObject again
(cacheOutput=true). The DataFrame then cumulates the transformations of several Actions, but the exported
lineage still respects the Action boundaries: a column of an input DataObject is a leaf of the lineage of the
Action reading it. How the input DataObject got that column is described by the export of the Action which
wrote it. There is no need to turn caching off to export lineage.
cacheInput has no effect on the lineage at all, as it only materializes DataFrames in the exec-phase.
If an Action reads both the cached output of a previous Action and a DataObject that previous Action passed through unchanged, both inputs share the very same columns. Such a column is reported for both input DataObjects, as there is no way to tell which one it was read from - and both are true.
Limitations
Column lineage is analyzed for the Spark engine, the plain-Scala engine and the SQL engine (sdl-sql). They use the same format and follow the same rules, but they read the lineage from different places: for Spark it is read from the expression ids of the analyzed logical plan, while the plain-Scala engine evaluates expressions immediately and therefore records with every column it creates which columns its values come from. The SQL engine reads it from the SQL query it creates, with the lineage module of SQLGlot.
The analysis is best-effort: a column which can not be traced back to an input DataObject is reported as unresolved rather than with a wrong source. This is the case for
- transformations which break the DataFrame lineage, e.g. a custom transformer which reads data itself, or a DataFrame created from data inside a transformation
Columns which are used in a join, filter, group by or sort condition, or which a window function is
partitioned and ordered by, influence the output without being part of its value. They are not reported yet - OpenLineage calls this INDIRECT lineage and collects it
in the dataset list of the facet. An aggregation over all rows such as count(*) belongs there too, and is
reported without input columns until then.
For a typed Dataset transformation of the Spark engine, e.g. a map over a case class, the Scala function is
opaque. Every output column of such a transformation is therefore reported to depend on every column it reads.
User defined functions need no special handling, as a UDF reads the columns of its arguments like any other
expression.