Реалізація AI-трекінгу лінійки даних (Data Lineage)
Уявіть: ви змінюєте структуру таблиці у staging шарі, і через день дашборди відділу продажів показують нулі. Причина — пропущений reference у SQL-трансформації. Без автоматичної лінійки даних знайти її — години ручного перебору. Коли пайплайнів десятки, а дашбордів сотні, ручний impact analysis стає вузьким горлечком. Наша AI-система лінійки даних автоматизує побудову графа, який покаже кожне джерело та трансформацію, і зробить impact analysis за секунди.
Наша компанія — експерт з 7-річним досвідом у Data Engineering, реалізувала 25+ проектів з Data Lineage. Маємо сертифікацію OpenLineage та гарантію точності графа. Економія на підтримці даних — до $50 000 на рік для середнього проекту.
Ми інтегруємо AI-трекінг походження даних, який автоматично будує граф із вашого SQL, dbt-моделей та Python-коду. Без цього знайти джерело помилки в колонці — години ручного пошуку. Наше рішення будує граф на льоту, без ручної документації, і покриває 80-90% типових ETL-патернів. Решта 10% — runtime-трекінг через OpenLineage.
Автоматизація лінійки даних: чому це критично
Втрата лінійки — одна з частих причин простоїв у дата-платформах. Зміна однієї колонки в raw-шарі може зламати 20+ дашбордів у BI. Без автоматичного трекінгу лінійки impact analysis займає дні, а помилки виявляються постфактум. AI-трекінг дає вам:
- граф залежностей — оновлюється після кожного деплою
- impact analysis — за секунди: "що зламається, якщо я видалю таблицю X"
- аудит трансформацій — звідки дані прийшли і як змінилися
Ручний impact analysis займає 2-3 дні, AI-трекінг — менше 100 мс на графі з 1000 вузлів: у 50-100 разів швидше. Економія бюджету на підтримку даних сягає 40%, а час на аудит скорочується з 2 днів до 2 годин.
Що дає impact analysis?
Impact analysis — це відповідь на питання "що зламається, якщо я зміню цю таблицю". Без нього кожна зміна — ризик. AI-трекінг дозволяє за секунди отримати повний список зачеплених дашбордів, моделей та пайплайнів. Наприклад, видалення колонки client_id у staging шарі може зачепити 15 дашбордів, 3 dbt-моделі та 2 API-ендпоїнти. Impact analysis покаже це до деплою.
Як ми будуємо граф лінійки даних?
Парсинг SQL та dbt-моделей
Для вилучення лінійки з SQL використовуємо sqlglot для AST-аналізу та LLM (Claude 3.5, GPT-4o) для складних випадків — динамічний SQL, віконні функції, ORM. Для dbt — читаємо manifest.json, будуємо граф на основі ref-залежностей.
from anthropic import Anthropic
import sqlparse
import sqlglot
import networkx as nx
import json
from dataclasses import dataclass
@dataclass
class LineageNode:
node_id: str
name: str
node_type: str # table, view, query, model, api, file
schema: dict = None
metadata: dict = None
@dataclass
class LineageEdge:
source: str
target: str
transform_type: str # select, join, aggregation, filter, union
columns_mapped: dict = None # {source_col: target_col}
class DataLineageTracker:
def __init__(self):
self.llm = Anthropic()
self.graph = nx.DiGraph()
self.nodes = {}
def parse_sql_lineage(self, sql: str, output_table: str = None) -> dict:
"""Извлечение линейки из SQL запроса"""
try:
# Парсинг через sqlglot
statements = sqlglot.parse(sql)
lineage = {'sources': [], 'targets': [], 'columns': {}}
for stmt in statements:
# Таблицы в FROM и JOIN
for table in stmt.find_all(sqlglot.expressions.Table):
if table.name:
lineage['sources'].append(table.name)
# Целевая таблица (CREATE TABLE AS / INSERT INTO)
if isinstance(stmt, sqlglot.expressions.Create):
lineage['targets'].append(str(stmt.this))
elif isinstance(stmt, sqlglot.expressions.Insert):
lineage['targets'].append(str(stmt.this))
if output_table:
lineage['targets'].append(output_table)
# Маппинг колонок через LLM для сложных случаев
lineage['column_mapping'] = self._extract_column_mapping(sql)
return lineage
except Exception:
# Fallback: LLM-парсинг
return self._llm_parse_lineage(sql, output_table)
def _llm_parse_lineage(self, sql: str, output_table: str = None) -> dict:
"""LLM-извлечение линейки для сложного SQL"""
response = self.llm.messages.create(
model="claude-3-5-sonnet-20241022",
max_tokens=400,
messages=[{
"role": "user",
"content": f"""Extract data lineage from this SQL.
SQL:
{sql[:1500]}
Output table: {output_table or "unknown"}
Return JSON:
{{
"sources": ["table1", "table2"],
"targets": ["output_table"],
"transforms": ["aggregation", "join"],
"column_mapping": {{"source.col1": "target.col_a"}}
}}"""
}]
)
try:
return json.loads(response.content[0].text)
except Exception:
return {'sources': [], 'targets': [], 'transforms': [], 'column_mapping': {}}
def _extract_column_mapping(self, sql: str) -> dict:
"""Маппинг колонок источник → цель"""
try:
parsed = sqlglot.parse_one(sql)
mapping = {}
for col in parsed.find_all(sqlglot.expressions.Column):
alias = col.find_ancestor(sqlglot.expressions.Alias)
if alias:
target_name = str(alias.alias)
source_name = str(col)
mapping[source_name] = target_name
return mapping
except Exception:
return {}
def add_dbt_lineage(self, manifest_path: str):
"""Импорт линейки из dbt manifest.json"""
with open(manifest_path) as f:
manifest = json.load(f)
for node_id, node in manifest.get('nodes', {}).items():
if node.get('resource_type') == 'model':
model_name = node['name']
# Добавление узла модели
self.graph.add_node(model_name, **{
'type': 'dbt_model',
'schema': node.get('database', '') + '.' + node.get('schema', ''),
'description': node.get('description', ''),
'tags': node.get('tags', [])
})
# Зависимости (upstream)
for dep in node.get('depends_on', {}).get('nodes', []):
dep_name = dep.split('.')[-1]
self.graph.add_edge(dep_name, model_name,
transform_type='dbt_ref')
def build_lineage_from_code(self, python_code: str,
file_name: str = "transform.py") -> dict:
"""Извлечение линейки из Python кода трансформации"""
response = self.llm.messages.create(
model="claude-3-5-sonnet-20241022",
max_tokens=400,
messages=[{
"role": "user",
"content": f"""Extract data lineage from this Python ETL code.
File: {file_name}
Code:
{python_code[:1500]}
Return JSON:
{{
"reads_from": ["table/file/api names"],
"writes_to": ["output table/file names"],
"transforms": ["description of transformations applied"],
"column_transforms": ["human-readable descriptions of column transformations"]
}}"""
}]
)
try:
return json.loads(response.content[0].text)
except Exception:
return {}
Граф лінійки та impact analysis
def get_downstream_impact(self, source_table: str) -> dict:
"""Что сломается если изменить source_table"""
if source_table not in self.graph:
return {'affected': [], 'count': 0}
# BFS для получения всех downstream узлов
affected = []
visited = set()
queue = [source_table]
while queue:
current = queue.pop(0)
if current in visited:
continue
visited.add(current)
successors = list(self.graph.successors(current))
for succ in successors:
if succ != source_table:
affected.append({
'node': succ,
'distance': nx.shortest_path_length(self.graph, source_table, succ),
'type': self.graph.nodes[succ].get('type', 'unknown')
})
queue.extend(successors)
affected.sort(key=lambda x: x['distance'])
return {
'source': source_table,
'affected': affected,
'count': len(affected)
}
def get_upstream_sources(self, target_table: str) -> dict:
"""Откуда данные в target_table"""
if target_table not in self.graph:
return {'sources': [], 'path': []}
# Все предки
ancestors = list(nx.ancestors(self.graph, target_table))
paths = {}
for source in ancestors:
try:
path = nx.shortest_path(self.graph, source, target_table)
paths[source] = path
except nx.NetworkXNoPath:
pass
# AI-объяснение линейки
explanation = self._explain_lineage(target_table, ancestors, paths)
return {
'target': target_table,
'sources': ancestors,
'paths': paths,
'explanation': explanation
}
def _explain_lineage(self, target: str, sources: list, paths: dict) -> str:
"""LLM-объяснение линейки данных"""
paths_summary = json.dumps(
{s: p for s, p in list(paths.items())[:5]},
ensure_ascii=False
)
response = self.llm.messages.create(
model="claude-3-5-sonnet-20241022",
max_tokens=300,
messages=[{
"role": "user",
"content": f"""Explain the data lineage for table "{target}".
Source tables: {sources}
Key paths: {paths_summary}
Summarize: where data originates, what transformations occur, potential data quality risks.
2-4 sentences, non-technical language."""
}]
)
return response.content[0].text
def detect_lineage_breaks(self) -> list[dict]:
"""Обнаружение разрывов в линейке"""
breaks = []
# Таблицы без источников (кроме raw)
for node in self.graph.nodes():
if self.graph.in_degree(node) == 0:
node_data = self.graph.nodes[node]
if node_data.get('type') not in ['raw_table', 'external_source']:
breaks.append({
'node': node,
'issue': 'no_upstream_lineage',
'severity': 'warning'
})
# Таблицы без потребителей
if self.graph.out_degree(node) == 0:
breaks.append({
'node': node,
'issue': 'orphaned_dataset',
'severity': 'info'
})
return breaks
Трекінг лінійки через парсинг SQL + доповнення LLM покриває 80-90% типових ETL-патернів. Для складних tricky трансформацій (динамічний SQL, ORM) потрібен runtime tracing. OpenLineage + Marquez — стандарт для автоматичного збору лінійки з Airflow, Spark та dbt без написання коду. Ми інтегруємо OpenLineage у ваш існуючий пайплайн за 1-2 дні, додаючи прозорість без зміни коду.
Покрокова інструкція впровадження AI-трекінгу
- Аудит — збираємо SQL, dbt, Python-код усіх пайплайнів.
- Парсинг — sqlglot + LLM вилучають залежності та мапінг колонок.
- Побудова графа — NetworkX створює DAG з вузлами та ребрами.
- Impact analysis — реалізуємо API для запитів downstream/upstream.
- Інтеграція в CI/CD — автоматичне оновлення графа при кожному деплої.
- Тестування — перевірка покриття та точності на реальних даних.
Кейс: автоматизація лінійки для рітейл-дашборду
У великому рітейлері 50+ дашбордів у Mode Analytics будувалися на 30 dbt-моделях. Після зміни структури у staging шарі (додавання колонки discount) один дашборд почав показувати некоректні суми. Без лінійки пошук причини зайняв 3 дні. Після впровадження AI-трекінгу impact analysis показав, що зачеплені 5 дашбордів та 2 dbt-моделі. Рішення: автоматичне оновлення графа при кожному коміті. Час на пошук помилок скоротився до 10 хвилин.
Що входить у результат роботи
- Автоматично оновлюваний граф лінійки даних (дашборд на основі NetworkX + D3.js).
- API для impact analysis з підтримкою upstream/downstream запитів.
- Детекція розривів у лінійці (таблиці без джерел або споживачів).
- Інтеграція з CI/CD (GitHub Actions / GitLab CI) для автооновлення.
- Документація архітектури та інструкції для команди.
- Навчання команди (2 години воркшоп).
- Підтримка 3 місяці після впровадження.
Етапи реалізації AI-трекінгу
| Етап | Зміст | Строк (днів) |
|---|---|---|
| Аудит поточних пайплайнів | Збір SQL, dbt, Python-коду; виявлення прогалин | 1-2 |
| Побудова графа | Парсинг, LLM-доповнення, ручне тестування | 3-7 |
| Impact analysis | API, дашборд, детекція розривів | 2-3 |
| Інтеграція в CI/CD | Автоматичне оновлення при деплої | 1-2 |
| Документація та навчання | README, дашборди, інструкція для команди | 1-2 |
Порівняння методів: ручний vs AI-трекінг
| Критерій | Ручний трекінг | AI-трекінг |
|---|---|---|
| Час на один пайплайн | 2-3 тижні | від 5 днів |
| Покриття автоматизацією | 0% | 80-90% |
| Latency impact analysis (1000 вузлів) | години | < 100 мс |
| Оновлення при деплої | вручну | автоматично |
AI-трекінг у 50 разів швидший за ручний impact analysis, а економія на підтримці даних перевищує 40%. Оцінимо ваш проект — пишіть. Наш досвід: 7+ років у Data Engineering, понад 25 реалізованих проектів з Data Lineage. Сертифіковані за OpenLineage. Отримайте консультацію — зв'яжіться з нами.







