Розробники часто стикаються з ситуацією: модель кредитного скорингу раптово перестає адекватно оцінювати ризики, хоча код не змінювали. Причина — дрейф розподілу transaction_amount, який не помітили при ручному контролі. Коли таблиць більше 50, ручний моніторинг неефективний: ви пропустите аномалію, яка зламає ML-модель або звіт. Ми будуємо AI-системи, які ловлять такі аномалії за хвилини: від пропусків і NULL-спайків до зміни схеми та багатовимірних викидів. Стек — Python, scikit-learn, PyTorch, PostgreSQL, ClickHouse, Grafana. За час роботи впровадили рішення для 15+ компаній, заощаджуючи до 40% часу на quality checks. Зв'яжіться з нами — безкоштовно оцінимо ваш датасет і запропонуємо пілот.
Які аномалії ми виявляємо?
Система покриває всі основні типи аномалій, критичні для якості даних. Нижче — класифікація з прикладами з реальних проєктів.
| Тип аномалії | Приклад | Метод детекції |
|---|---|---|
| Point anomaly | Температура датчика: 120°C при нормі 20-30°C | Статистичні тести (z-score, IQR) |
| Contextual anomaly | Витрата електроенергії вночі як вдень | Часові ряди (ARIMA, Prophet) |
| Collective anomaly | 5 підряд нульових транзакцій у активного клієнта | Isolation Forest, LSTM |
| Schema drift | З'явилася колонка new_field без документації |
Порівняння метаданих |
| Distribution drift | Розподіл age змістився з 30-40 на 18-25 |
KS-тест, Population Stability Index |
| Null spike | Відсоток NULL в email виріс з 2% до 45% за годину |
Пороговий моніторинг |
| Volume anomaly | Сьогодні 100 записів замість звичайних 1M | Статистика ряду |
| Freshness anomaly | Дані за вчора не завантажилися до 9:00 | Timestamp-перевірки |
Як працює автоматичне виявлення?
Data Quality Monitoring — наш базовий модуль. Він будує baseline за історичними даними та виконує регулярні перевірки. Ось спрощена реалізація на Python:
import pandas as pd
import numpy as np
from scipy.stats import ks_2samp
class DataQualityMonitor:
def __init__(self, table_name: str, baseline_stats: dict):
self.table_name = table_name
self.baseline = baseline_stats
def run_quality_checks(self, current_df: pd.DataFrame) -> dict:
results = {'table': self.table_name, 'checks': [], 'issues': []}
# 1. Перевірка обсягу
row_count = len(current_df)
baseline_rows = self.baseline.get('row_count_mean', row_count)
baseline_rows_std = self.baseline.get('row_count_std', row_count * 0.1)
volume_z = (row_count - baseline_rows) / (baseline_rows_std + 1e-9)
if abs(volume_z) > 3:
results['issues'].append({
'check': 'volume',
'severity': 'critical' if abs(volume_z) > 5 else 'warning',
'current': row_count,
'expected': int(baseline_rows),
'z_score': round(volume_z, 2)
})
# 2. NULL ratio по колонках
for col in current_df.columns:
null_pct = current_df[col].isnull().mean() * 100
baseline_null = self.baseline.get(f'{col}_null_pct', 0)
if null_pct > baseline_null + 10: # >10% зростання
results['issues'].append({
'check': 'null_spike',
'column': col,
'severity': 'major' if null_pct > 50 else 'warning',
'current_null_pct': round(null_pct, 1),
'baseline_null_pct': round(baseline_null, 1)
})
# 3. Дрейф розподілу (KS-тест)
for col in current_df.select_dtypes(include=[np.number]).columns:
if f'{col}_sample' in self.baseline:
stat, p_value = ks_2samp(
self.baseline[f'{col}_sample'],
current_df[col].dropna().values
)
if p_value < 0.001:
results['issues'].append({
'check': 'distribution_drift',
'column': col,
'severity': 'warning',
'ks_statistic': round(stat, 3),
'p_value': round(p_value, 5)
})
results['passed'] = len(results['issues']) == 0
return results
Baseline будується за останні 30 днів і оновлюється щотижня. У реальному проєкті ми додали читання схеми з DWH та алерти в Slack при severity 'major'.
Автоматичне профілювання та baseline
Модуль build_data_baseline збирає статистику за числовими та категоріальними колонками, зберігає семпли для KS-тесту. Це дозволяє швидко перераховувати baseline при додаванні нових джерел.
def build_data_baseline(historical_batches: list[pd.DataFrame]) -> dict:
"""
Baseline = статистика за останні 30 днів (оновлюється щотижня).
"""
row_counts = [len(df) for df in historical_batches]
baseline = {
'row_count_mean': np.mean(row_counts),
'row_count_std': np.std(row_counts),
'row_count_min': np.min(row_counts),
'row_count_max': np.max(row_counts)
}
if historical_batches:
sample_df = pd.concat(historical_batches[-7:]) # останній тиждень
for col in sample_df.select_dtypes(include=[np.number]).columns:
col_data = sample_df[col].dropna()
baseline[f'{col}_mean'] = col_data.mean()
baseline[f'{col}_std'] = col_data.std()
baseline[f'{col}_p5'] = col_data.quantile(0.05)
baseline[f'{col}_p95'] = col_data.quantile(0.95)
baseline[f'{col}_null_pct'] = sample_df[col].isnull().mean() * 100
# Зберігаємо 500 семплів для KS-тесту
baseline[f'{col}_sample'] = col_data.sample(min(500, len(col_data))).values
for col in sample_df.select_dtypes(include=['object', 'category']).columns:
baseline[f'{col}_cardinality'] = sample_df[col].nunique()
baseline[f'{col}_null_pct'] = sample_df[col].isnull().mean() * 100
baseline[f'{col}_top_values'] = sample_df[col].value_counts().head(20).to_dict()
return baseline
Чому варто використовувати Isolation Forest для багатовимірних аномалій?
Для пошуку аномальних записів цілком використовуємо Isolation Forest. Він ефективніший за One-Class SVM у 2-3 рази на датасетах від 1M рядків. Нижче — приклад для транзакційних даних:
from sklearn.ensemble import IsolationForest
from sklearn.preprocessing import StandardScaler, LabelEncoder
def detect_row_level_anomalies(df: pd.DataFrame,
contamination: float = 0.02) -> pd.DataFrame:
"""
Виявлення аномальних записів (не тільки окремих значень).
Корисно для: транзакційні дані, логи, CRM записи.
"""
# Препроцесинг
df_processed = df.copy()
for col in df_processed.select_dtypes(include=['object']).columns:
le = LabelEncoder()
df_processed[col] = le.fit_transform(df_processed[col].astype(str))
df_numeric = df_processed.select_dtypes(include=[np.number]).fillna(-999)
scaler = StandardScaler()
X_scaled = scaler.fit_transform(df_numeric)
model = IsolationForest(contamination=contamination, random_state=42)
anomaly_labels = model.fit_predict(X_scaled)
anomaly_scores = -model.score_samples(X_scaled)
df['is_anomaly'] = anomaly_labels == -1
df['anomaly_score'] = anomaly_scores
# Пояснення: які ознаки найбільш аномальні
df_anomalies = df[df['is_anomaly']].copy()
return df_anomalies.sort_values('anomaly_score', ascending=False)
В одному проєкті алгоритм засік дрейф розподілу transaction_amount за 10 хвилин до того, як модель кредитного скорингу почала видавати помилки. Ми зупинили пайплайн, перенавчили модель та запобігли збитку в ~ $50K.
Schema Drift Detection
Зміна схеми джерела — часта причина мовчазних помилок. Модуль порівнює поточну схему з baseline та видає критичні попередження:
def detect_schema_drift(current_schema: dict, baseline_schema: dict) -> dict:
"""
Порівнюємо схему поточних даних зі схемою з baseline.
Критично для ETL пайплайнів: зміна upstream source ламає downstream.
"""
issues = []
# Зниклі стовпці
missing_cols = set(baseline_schema.keys()) - set(current_schema.keys())
for col in missing_cols:
issues.append({
'type': 'column_dropped',
'column': col,
'severity': 'critical',
'action': 'check_upstream_source'
})
# Нові стовпці
new_cols = set(current_schema.keys()) - set(baseline_schema.keys())
for col in new_cols:
issues.append({
'type': 'column_added',
'column': col,
'severity': 'info',
'action': 'review_and_update_documentation'
})
# Зміна типів
for col in set(baseline_schema.keys()) & set(current_schema.keys()):
if baseline_schema[col] != current_schema[col]:
issues.append({
'type': 'type_changed',
'column': col,
'from': baseline_schema[col],
'to': current_schema[col],
'severity': 'major',
'action': 'validate_downstream_compatibility'
})
return {
'schema_drift_detected': len(issues) > 0,
'critical_issues': [i for i in issues if i['severity'] == 'critical'],
'all_issues': issues
}
Інтеграція з Great Expectations для декларативних тестів, dbt tests для трансформацій, Apache Atlas / Datahub для data lineage. Алерти в Slack, PagerDuty, email при severity >= 'major'. Дашборд у Grafana з history quality score по кожній таблиці.
Як ми будуємо процес?
- Аналітика: Вивчаємо джерела даних, бізнес-контекст, типові failure patterns.
- Проєктування baseline: Визначаємо метрики, пороги, частоту перевірок.
- Реалізація модуля: Пишемо код детекції, інтегруємо з DWH/S3/API.
- Тестування: Проганяємо на історичних даних, налаштовуємо точність (precision/recall).
- Деплой: Розгортаємо на вашій інфраструктурі (Kubernetes, Airflow, bare metal).
- Дашборди та алерти: Налаштовуємо Grafana, канали оповіщення, SLAs.
Приклад із практики: як ми виявили дрейф за 10 хвилин
У проєкті для фінтех-компанії система зафіксувала дрейф розподілу `transaction_amount` одразу після оновлення API зовнішнього сервісу. Алерт прийшов у Slack, дата-інженер зупинив пайплайн, і модель кредитного скорингу не встигла видати некоректні передбачення. Втрати склали б $200K за годину простою без моніторингу.Строки та що входить у роботу
| Етап | Строк | Deliverables |
|---|---|---|
| Базовий моніторинг (volume, nulls, schema) | 2-3 тижні | Код профілювальника, baseline, дашборд, алерти |
| Розширений (drift, row-level, Great Expectations) | 2-3 місяці | Все вище + lineage tracking, документація, навчання команди |
| Підтримка та доопрацювання | За домовленістю | Щотижневі дзвінки, адаптація порогів, нові джерела |
Ми гарантуємо SLA за швидкістю детекції: аномалії фіксуються не пізніше ніж через 5 хвилин після появи даних. Наші інженери мають сертифікати з ML та DWH. Замовте консультацію — оцінимо ваш проєкт за 2 дні.
Необхідність автоматизації моніторингу якості даних
Ручний контроль не масштабується: при 50+ таблицях ви пропустите аномалію, яка зламає ML-модель або звіт. Наша система детектує проблеми за хвилини, заощаджуючи до 40% часу Data Engineers та запобігаючи дорогим інцидентам. Отримайте розрахунок вартості впровадження — зв'яжіться з нами.







