Skip to content

Commit 96fb5f2

Browse files
kalra-mohitclaude
andcommitted
Fix column lineage transformation details always appearing as null (#3100)
Marquez's internal LineageEvent.ColumnLineageInputField model (used to deserialize the OpenLineage ColumnLineageDatasetFacet on ingestion) only captured namespace/name/field for each input field - it had no `transformations` property at all. The write path (OpenLineageDao.upsertColumnLineage) only ever read transformationDescription/transformationType off the *output column* object (ColumnLineageOutputColumn), which per the current OpenLineage spec are explicitly deprecated: https://github.com/OpenLineage/OpenLineage/blob/main/spec/facets/ColumnLineageDatasetFacet.json Modern producers report transformation details per input field instead, via a `transformations` array nested under each entry of `inputFields[]` (each with type/subtype/description/masking) - this is exactly the format shown in the issue's example JSON. Since Marquez's ingestion model didn't have a `transformations` field to deserialize that array into, and never read the deprecated output-column-level fields from producers that only emit the new format, transformation_description/transformation_type were written as null for any producer using the current (non-deprecated) facet shape - reproducing the bug exactly as reported. Fix: - Added `transformations` (List<Transformation>) to LineageEvent.ColumnLineageInputField, with a new nested Transformation class (type, subtype, description, masking) mirroring the OpenLineage spec. The existing 3-arg constructor (namespace, name, field) is preserved unchanged, so no existing call site needed to be touched. - Added OpenLineageDao.transformationOf(inputField, outputColumn), a small pure function that resolves the (description, type) pair for a given input field: prefer the first entry of that specific input field's (non-deprecated) `transformations` list; fall back to the deprecated whole-output-column transformationDescription/transformationType for producers that still only report the old format. - Updated OpenLineageDao.upsertColumnLineage to resolve the transformation per matched input field (via transformationOf) and group input fields by their resolved (description, type) before calling the existing ColumnLineageDao.upsertColumnLineageRow(...) batch API once per group. This correctly persists a distinct transformation per input/output field edge - which the column_lineage table and the read path (ColumnLineageService, which already surfaces per-input-field transformation info) both already supported - without changing ColumnLineageDao's public method signature, so no existing ColumnLineageDaoTest call sites needed to change. Tests: - api/src/test/java/marquez/db/OpenLineageDaoColumnLineageTransformationTest.java (new, tagged UnitTests, no DB required): exercises OpenLineageDao.transformationOf(...) directly against the exact scenario from the issue (an output column with no deprecated fields set, fed by input fields that each report their own `transformations` entry), plus the multi-input-field, fallback, "both formats absent", and "multiple transformations reported" cases. I verified these tests explicitly: temporarily reverted transformationOf(...) to the old behavior (always read the deprecated output-column-level fields) and confirmed 3 of 5 tests fail with the expected description/type coming back null - then restored the fix and confirmed all 5 pass. - OpenLineageDaoTest#testUpdateMarquezModelDatasetWithColumnLineageFacet_usesNewTransformationsFormat (new): full DB-backed regression test that builds a ColumnLineageDatasetFacet using only the new inputFields[].transformations[] format (deprecated fields left unset, as real modern producers do) and asserts the persisted ColumnLineageRow has the expected non-null transformation description/type. Test evidence: ./gradlew :api:testUnit passes (123 tests, including the 5 new pure-function tests). I was not able to execute OpenLineageDaoTest (Postgres-backed) locally in this sandbox - Testcontainers fails to negotiate with the local Docker Engine (old docker-java client defaults to Docker API v1.32, local engine requires >= v1.40); this is a pre-existing environment limitation that affects every DB-backed test in the suite, not something introduced by this change (verified identically on unmodified tests like RunDaoTest). ./gradlew :api:compileTestJava passes, confirming the new integration test at least compiles correctly against the real DAO/model classes. Fixes #3100 Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Signed-off-by: Mohit Kalra <mohit2494@gmail.com>
1 parent 180f37b commit 96fb5f2

4 files changed

Lines changed: 339 additions & 24 deletions

File tree

‎api/src/main/java/marquez/db/OpenLineageDao.java‎

Lines changed: 74 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -1049,10 +1049,23 @@ private List<ColumnLineageRow> upsertColumnLineage(
10491049
return Stream.empty();
10501050
}
10511051

1052-
// get field uuids of input columns related to this run
1053-
List<Pair<UUID, UUID>> inputFields =
1052+
// Match each field associated with this run to the OpenLineage input field it
1053+
// corresponds to (if any), along with the transformation description/type that
1054+
// applies to that specific input -> output edge. Prefer the (non-deprecated)
1055+
// per-input-field InputField#transformations entry reported by the producer;
1056+
// fall back to the deprecated, whole-output-column
1057+
// transformationDescription/transformationType for producers that only report the
1058+
// old format. Without this fallback/lookup, transformation details reported via the
1059+
// modern `transformations` array were silently dropped and always appeared as null
1060+
// (see #3100).
1061+
//
1062+
// Input fields are then grouped by the transformation (description, type) that
1063+
// applies to each of them, so that each distinct transformation is persisted with
1064+
// its own value while still reusing the existing batch upsertColumnLineageRow(...)
1065+
// API (which applies a single description/type to every input field passed to it).
1066+
Map<Pair<String, String>, List<Pair<UUID, UUID>>> inputFieldsByTransformation =
10541067
runFields.stream()
1055-
.filter(
1068+
.flatMap(
10561069
fieldData ->
10571070
columnLineage.getInputFields().stream()
10581071
.filter(
@@ -1061,33 +1074,72 @@ private List<ColumnLineageRow> upsertColumnLineage(
10611074
&& of.getName().equals(fieldData.getDatasetName())
10621075
&& of.getField().equals(fieldData.getField()))
10631076
.findAny()
1064-
.isPresent())
1065-
.map(
1066-
fieldData ->
1067-
Pair.of(
1068-
fieldData.getDatasetVersionUuid(),
1069-
fieldData.getDatasetFieldUuid()))
1070-
.collect(Collectors.toList());
1077+
.map(matchedInputField -> Pair.of(fieldData, matchedInputField))
1078+
.stream())
1079+
.collect(
1080+
Collectors.groupingBy(
1081+
fieldDataAndInputField ->
1082+
transformationOf(
1083+
fieldDataAndInputField.getRight(), columnLineage),
1084+
Collectors.mapping(
1085+
fieldDataAndInputField ->
1086+
Pair.of(
1087+
fieldDataAndInputField.getLeft().getDatasetVersionUuid(),
1088+
fieldDataAndInputField.getLeft().getDatasetFieldUuid()),
1089+
Collectors.toList())));
10711090

10721091
log.debug(
1073-
"Adding column lineage on output field '{}' for dataset version '{}' with input fields: {}",
1092+
"Adding column lineage on output field '{}' for dataset version '{}' with input fields by transformation: {}",
10741093
outputField.get().getName(),
10751094
datasetVersionRow.getUuid(),
1076-
inputFields);
1077-
return daos
1078-
.getColumnLineageDao()
1079-
.upsertColumnLineageRow(
1080-
datasetVersionRow.getUuid(),
1081-
outputField.get().getUuid(),
1082-
inputFields,
1083-
columnLineage.getTransformationDescription(),
1084-
columnLineage.getTransformationType(),
1085-
now)
1086-
.stream();
1095+
inputFieldsByTransformation);
1096+
1097+
return inputFieldsByTransformation.entrySet().stream()
1098+
.flatMap(
1099+
entry ->
1100+
daos.getColumnLineageDao()
1101+
.upsertColumnLineageRow(
1102+
datasetVersionRow.getUuid(),
1103+
outputField.get().getUuid(),
1104+
entry.getValue(),
1105+
entry.getKey().getLeft(),
1106+
entry.getKey().getRight(),
1107+
now)
1108+
.stream());
10871109
})
10881110
.collect(Collectors.toList());
10891111
}
10901112

1113+
/**
1114+
* Resolves the (transformationDescription, transformationType) pair that applies to a single
1115+
* OpenLineage column-lineage input field. Prefers the first entry of the (non-deprecated) {@link
1116+
* LineageEvent.ColumnLineageInputField#getTransformations()} list reported for that specific
1117+
* input field; falls back to the deprecated, whole-output-column {@code
1118+
* transformationDescription}/{@code transformationType} on {@link
1119+
* LineageEvent.ColumnLineageOutputColumn} for producers that only report the old format.
1120+
*
1121+
* <p>If a producer reports more than one transformation for a single input field, only the
1122+
* first is used, matching the granularity at which Marquez persists a transformation
1123+
* description/type (one value per input/output field edge).
1124+
*/
1125+
static Pair<String, String> transformationOf(
1126+
LineageEvent.ColumnLineageInputField inputField,
1127+
LineageEvent.ColumnLineageOutputColumn outputColumn) {
1128+
Optional<LineageEvent.ColumnLineageInputField.Transformation> transformation =
1129+
Optional.ofNullable(inputField.getTransformations()).stream()
1130+
.flatMap(List::stream)
1131+
.findFirst();
1132+
String transformationDescription =
1133+
transformation
1134+
.map(LineageEvent.ColumnLineageInputField.Transformation::getDescription)
1135+
.orElseGet(outputColumn::getTransformationDescription);
1136+
String transformationType =
1137+
transformation
1138+
.map(LineageEvent.ColumnLineageInputField.Transformation::getType)
1139+
.orElseGet(outputColumn::getTransformationType);
1140+
return Pair.of(transformationDescription, transformationType);
1141+
}
1142+
10911143
default String formatDatasetName(String name) {
10921144
return name;
10931145
}

‎api/src/main/java/marquez/service/models/LineageEvent.java‎

Lines changed: 38 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,8 +630,6 @@ public static class ColumnLineageOutputColumn extends BaseJsonModel {
630630
private String transformationType;
631631
}
632632

633-
@Builder
634-
@AllArgsConstructor
635633
@NoArgsConstructor
636634
@Setter
637635
@Getter
@@ -642,6 +640,44 @@ public static class ColumnLineageInputField extends BaseJsonModel {
642640
@NotNull private String namespace;
643641
@NotNull private String name;
644642
@NotNull private String field;
643+
644+
/**
645+
* The transformations applied to this specific input field in order to compute the output
646+
* field it is associated with. This is the current (non-deprecated) way for a producer to
647+
* report per-input-field transformation details, as opposed to the deprecated {@link
648+
* ColumnLineageOutputColumn#getTransformationDescription()} / {@link
649+
* ColumnLineageOutputColumn#getTransformationType()}, which apply (at most) a single,
650+
* whole-column value shared by every input field. See
651+
* https://github.com/OpenLineage/OpenLineage/blob/main/spec/facets/ColumnLineageDatasetFacet.json
652+
*/
653+
private List<Transformation> transformations;
654+
655+
public ColumnLineageInputField(String namespace, String name, String field) {
656+
this(namespace, name, field, null);
657+
}
658+
659+
@Builder
660+
public ColumnLineageInputField(
661+
String namespace, String name, String field, List<Transformation> transformations) {
662+
this.namespace = namespace;
663+
this.name = name;
664+
this.field = field;
665+
this.transformations = transformations;
666+
}
667+
668+
@Builder
669+
@AllArgsConstructor
670+
@NoArgsConstructor
671+
@Setter
672+
@Getter
673+
@Valid
674+
@ToString
675+
public static class Transformation {
676+
private String type;
677+
private String subtype;
678+
private String description;
679+
private Boolean masking;
680+
}
645681
}
646682

647683
@Builder
Lines changed: 153 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,153 @@
1+
/*
2+
* Copyright 2018-2023 contributors to the Marquez project
3+
* SPDX-License-Identifier: Apache-2.0
4+
*/
5+
6+
package marquez.db;
7+
8+
import static org.assertj.core.api.Assertions.assertThat;
9+
10+
import java.util.Arrays;
11+
import java.util.Collections;
12+
import java.util.List;
13+
import marquez.service.models.LineageEvent;
14+
import marquez.service.models.LineageEvent.ColumnLineageInputField;
15+
import marquez.service.models.LineageEvent.ColumnLineageInputField.Transformation;
16+
import marquez.service.models.LineageEvent.ColumnLineageOutputColumn;
17+
import org.apache.commons.lang3.tuple.Pair;
18+
import org.junit.jupiter.api.Test;
19+
20+
/**
21+
* Regression tests for https://github.com/MarquezProject/marquez/issues/3100: transformation
22+
* details reported via the (non-deprecated) {@code InputField.transformations[]} array of the
23+
* OpenLineage {@code ColumnLineageDatasetFacet} were always dropped, so {@code
24+
* transformation_description}/{@code transformation_type} always ended up null in Marquez,
25+
* regardless of what the producer actually reported.
26+
*
27+
* <p>These tests exercise {@link OpenLineageDao#transformationOf(ColumnLineageInputField,
28+
* ColumnLineageOutputColumn)} directly - the pure function responsible for resolving which
29+
* transformation applies to a given input field - so they can run without a database.
30+
*/
31+
@org.junit.jupiter.api.Tag("UnitTests")
32+
class OpenLineageDaoColumnLineageTransformationTest {
33+
34+
@Test
35+
void prefersPerInputFieldTransformationOverDeprecatedOutputColumnFields() {
36+
// Reproduces the exact scenario from the issue: the output column itself does not set the
37+
// deprecated transformationDescription/transformationType, but each input field reports its
38+
// own transformations[] entry.
39+
ColumnLineageInputField inputField =
40+
ColumnLineageInputField.builder()
41+
.namespace("bookstore2")
42+
.name("customers")
43+
.field("customer_email")
44+
.transformations(
45+
Collections.singletonList(
46+
Transformation.builder()
47+
.type("DIRECT")
48+
.subtype("TRANSFORMATION")
49+
.description("concat(customers.customer_name, ' - ', customers.customer_email)")
50+
.build()))
51+
.build();
52+
ColumnLineageOutputColumn outputColumn =
53+
ColumnLineageOutputColumn.builder().inputFields(List.of(inputField)).build();
54+
55+
Pair<String, String> transformation = OpenLineageDao.transformationOf(inputField, outputColumn);
56+
57+
assertThat(transformation.getLeft())
58+
.isEqualTo("concat(customers.customer_name, ' - ', customers.customer_email)");
59+
assertThat(transformation.getRight()).isEqualTo("DIRECT");
60+
}
61+
62+
@Test
63+
void distinctInputFieldsOnTheSameOutputColumnCanHaveDifferentTransformations() {
64+
// customer_full in the issue is fed by two input fields, each with a different underlying
65+
// source field but (in the issue's example) the same description; verify the resolution is
66+
// genuinely per-input-field rather than picking a single value for the whole output column.
67+
ColumnLineageInputField emailField =
68+
ColumnLineageInputField.builder()
69+
.namespace("bookstore2")
70+
.name("customers")
71+
.field("customer_email")
72+
.transformations(
73+
Collections.singletonList(
74+
Transformation.builder().type("DIRECT").description("descriptionA").build()))
75+
.build();
76+
ColumnLineageInputField nameField =
77+
ColumnLineageInputField.builder()
78+
.namespace("bookstore2")
79+
.name("customers")
80+
.field("customer_name")
81+
.transformations(
82+
Collections.singletonList(
83+
Transformation.builder().type("INDIRECT").description("descriptionB").build()))
84+
.build();
85+
ColumnLineageOutputColumn outputColumn =
86+
ColumnLineageOutputColumn.builder()
87+
.inputFields(Arrays.asList(emailField, nameField))
88+
.build();
89+
90+
Pair<String, String> emailTransformation =
91+
OpenLineageDao.transformationOf(emailField, outputColumn);
92+
Pair<String, String> nameTransformation =
93+
OpenLineageDao.transformationOf(nameField, outputColumn);
94+
95+
assertThat(emailTransformation).isEqualTo(Pair.of("descriptionA", "DIRECT"));
96+
assertThat(nameTransformation).isEqualTo(Pair.of("descriptionB", "INDIRECT"));
97+
}
98+
99+
@Test
100+
void fallsBackToDeprecatedOutputColumnFieldsWhenInputFieldHasNoTransformations() {
101+
// Producers using the old (deprecated) format never set inputField.transformations at all;
102+
// this must keep working exactly as before.
103+
ColumnLineageInputField inputField =
104+
ColumnLineageInputField.builder()
105+
.namespace("ns")
106+
.name("upstream")
107+
.field("col")
108+
.build(); // no transformations set
109+
ColumnLineageOutputColumn outputColumn =
110+
ColumnLineageOutputColumn.builder()
111+
.inputFields(List.of(inputField))
112+
.transformationDescription("legacy description")
113+
.transformationType("IDENTITY")
114+
.build();
115+
116+
Pair<String, String> transformation = OpenLineageDao.transformationOf(inputField, outputColumn);
117+
118+
assertThat(transformation).isEqualTo(Pair.of("legacy description", "IDENTITY"));
119+
}
120+
121+
@Test
122+
void returnsNullPairWhenNeitherFormatIsReported() {
123+
ColumnLineageInputField inputField =
124+
ColumnLineageInputField.builder().namespace("ns").name("upstream").field("col").build();
125+
ColumnLineageOutputColumn outputColumn =
126+
ColumnLineageOutputColumn.builder().inputFields(List.of(inputField)).build();
127+
128+
Pair<String, String> transformation = OpenLineageDao.transformationOf(inputField, outputColumn);
129+
130+
assertThat(transformation.getLeft()).isNull();
131+
assertThat(transformation.getRight()).isNull();
132+
}
133+
134+
@Test
135+
void usesOnlyFirstTransformationWhenMultipleAreReportedForTheSameInputField() {
136+
ColumnLineageInputField inputField =
137+
ColumnLineageInputField.builder()
138+
.namespace("ns")
139+
.name("upstream")
140+
.field("col")
141+
.transformations(
142+
Arrays.asList(
143+
Transformation.builder().type("DIRECT").description("first").build(),
144+
Transformation.builder().type("INDIRECT").description("second").build()))
145+
.build();
146+
ColumnLineageOutputColumn outputColumn =
147+
ColumnLineageOutputColumn.builder().inputFields(List.of(inputField)).build();
148+
149+
Pair<String, String> transformation = OpenLineageDao.transformationOf(inputField, outputColumn);
150+
151+
assertThat(transformation).isEqualTo(Pair.of("first", "DIRECT"));
152+
}
153+
}

0 commit comments

Comments
 (0)