Issue with column transformation
package com.lineage.service; import com.lineage.model.ColumnLineageEntry; import com.lineage.model.LineageResult; import io.openlineage.sql.ColumnLineage; import io.openlineage.sql.ColumnMeta; import io.openlineage.sql.DbTableMeta; import io.openlineage.sql.OpenLineageSql; import io.openlineage.sql.SqlMeta; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import java.util.*; import java.util.stream.Collectors; /** * Extracts column-level lineage from SQL using OpenLineage's Rust-backed SQL parser. * * Key OpenLineage SQL parser concepts: * - OpenLineageSql.parse(sqls, dialect) → Optional* - SqlMeta.inTables() → List (source tables) * - SqlMeta.outTables() → List (target tables) * - SqlMeta.columnLineage() → List (column-level mappings) * - ColumnLineage.descendant → ColumnMeta (target column) * - ColumnLineage.lineage → List (source columns) * - ColumnMeta.name() → String (column name) * - ColumnMeta.origin() → DbTableMeta (which table it comes from) * - DbTableMeta: database(), schema(), name() */ @Slf4j @Service public class SqlColumnLineageService { private static final String DEFAULT_DIALECT = "ansi"; /** * Parse one or more SQL statements and extract full lineage. * * @param sql the SQL string (can contain multiple statements separated by ';') * @param dialect SQL dialect (postgres, bigquery, snowflake, hive, sparksql, etc.) * @return LineageResult with table-level + column-level lineage */ public LineageResult extractLineage(String sql, String dialect) { String resolvedDialect = (dialect != null && !dialect.isBlank()) ? dialect.trim().toLowerCase() : DEFAULT_DIALECT; log.info("Parsing SQL with dialect='{}', length={}", resolvedDialect, sql.length()); try { Optional sqlMetaOpt = OpenLineageSql.parse( Collections.singletonList(sql), resolvedDialect ); if (sqlMetaOpt.isEmpty()) { log.warn("OpenLineageSql.parse returned empty for dialect={}", resolvedDialect); return LineageResult.builder() .sql(sql) .dialect(resolvedDialect) .errors(List.of("Parser returned no result. Check SQL syntax or dialect.")) .build(); } SqlMeta meta = sqlMetaOpt.get(); // ── Table-level lineage ────────────────────────────── List inputTables = toTableNames(meta.inTables()); List outputTables = toTableNames(meta.outTables()); // ── Column-level lineage ───────────────────────────── List columnEntries = extractColumnLineages(meta.columnLineage()); log.info("Extracted: inputs={}, outputs={}, columnMappings={}", inputTables.size(), outputTables.size(), columnEntries.size()); return LineageResult.builder() .sql(sql) .dialect(resolvedDialect) .inputTables(inputTables) .outputTables(outputTables) .columnLineages(columnEntries) .build(); } catch (Exception e) { log.error("Failed to parse SQL: {}", e.getMessage(), e); return LineageResult.builder() .sql(sql) .dialect(resolvedDialect) .errors(List.of("Parse error: " + e.getMessage())) .build(); } } /** * Batch parse: multiple SQL statements individually. */ public List extractLineageBatch(List sqls, String dialect) { return sqls.stream() .map(sql -> extractLineage(sql, dialect)) .collect(Collectors.toList()); } // ─── Private helpers ───────────────────────────────────────── private List extractColumnLineages(List columnLineages) { if (columnLineages == null || columnLineages.isEmpty()) { return Collections.emptyList(); } List entries = new ArrayList<>(); for (ColumnLineage cl : columnLineages) { ColumnMeta descendant = cl.descendant(); String targetColumn = descendant.name(); List sources = cl.lineage(); if (sources == null || sources.isEmpty()) { // Target column exists but no traceable source (literal, constant, expression) entries.add(ColumnLineageEntry.builder() .sourceTable("EXPRESSION/LITERAL") .sourceColumn("-") .transformation("DERIVED") .targetColumn(targetColumn) .build()); continue; } for (ColumnMeta source : sources) { DbTableMeta origin = source.origin(); ColumnLineageEntry.ColumnLineageEntryBuilder builder = ColumnLineageEntry.builder() .sourceColumn(source.name()) .targetColumn(targetColumn); if (origin != null) { builder.sourceDatabase(origin.database()) .sourceSchema(origin.schema()) .sourceTable(origin.name()); } else { builder.sourceTable("UNRESOLVED"); } // Determine transformation type String transformation = determineTransformation(source, descendant, sources.size()); builder.transformation(transformation); entries.add(builder.build()); } } return entries; } /** * Heuristic to classify transformation type. * OpenLineage's SQL parser doesn't always provide transformation metadata, * so we infer from the mapping structure. */ private String determineTransformation(ColumnMeta source, ColumnMeta descendant, int sourceCount) { if (sourceCount == 1 && source.name().equalsIgnoreCase(descendant.name())) { return "DIRECT"; } else if (sourceCount == 1) { return "RENAMED"; } else { return "DERIVED"; // multiple source cols → expression/aggregation } } /** * Convert DbTableMeta list to fully qualified table name strings. */ private List toTableNames(List tables) { if (tables == null) return Collections.emptyList(); return tables.stream() .map(this::toFqtn) .collect(Collectors.toList()); } /** * Fully Qualified Table Name: database.schema.table */ private String toFqtn(DbTableMeta table) { StringBuilder sb = new StringBuilder(); if (table.database() != null) sb.append(table.database()).append("."); if (table.schema() != null) sb.append(table.schema()).append("."); sb.append(table.name()); return sb.toString(); } }
SqlColumnLineageService.java — calls OpenLineageSql.parse(sql, dialect) which returns SqlMeta containing columnLineage(), inTables(), and outTables() out of the box. No manual AST walking needed.
mvn clean install, and hit POST /api/lineage/extract with your SQL. The test file covers CTEs, triple-nested subqueries, scalar subqueries, CTAS, Snowflake QUALIFY, and BigQuery syntax.