Перейти к содержанию

Архитектура

Техническая архитектура Kruma Quality Validator — компонентная модель, потоки данных, диаграммы классов и точки расширения.

Версия: 1.0.1 | Владелец: Kruma Platform Team


Обзор системы

Quality Validator — Python-фреймворк для высокопроизводительной потоковой валидации данных. Проверяет записи на соответствие правилам из контрактов и маршрутизирует их в production или DLQ.

Высокоуровневая архитектура

                                KRUMA QUALITY VALIDATOR
    ┌────────────────────────────────────────────────────────────────┐
    │                                                                │
    │  ┌────────────────┐  ┌──────────┐  ┌──────────────────────┐   │
    │  │ ContractLoader │─▶│ Contract │─▶│  QualityValidator    │   │
    │  └────────────────┘  └──────────┘  └──────────┬───────────┘   │
    │         │                  │                   │               │
    │         ▼                  ▼                   ▼               │
    │  ┌────────────────┐  ┌──────────┐  ┌──────────────────────┐   │
    │  │  YAML-файлы    │  │  Модели  │  │     RuleEngine       │   │
    │  │  contract.yaml │  │  правил  │  │  Валидация + batch   │   │
    │  └────────────────┘  └──────────┘  └──────────────────────┘   │
    └────────────────────────────────────────────────────────────────┘
    ┌────────────────────────────────────────────────────────────────┐
    │                     KAFKA-ИНТЕГРАЦИЯ                          │
    │                                                                │
    │  ┌──────────────┐  ┌──────────────────┐  ┌──────────────┐    │
    │  │ Kafka        │─▶│ QualityValidator │─▶│ Kafka        │    │
    │  │ Consumer     │  │ Consumer         │  │ Producer     │    │
    │  │ (.raw)       │  │                  │  │ (.prod/.dlq) │    │
    │  └──────────────┘  └──────────────────┘  └──────────────┘    │
    └────────────────────────────────────────────────────────────────┘

Принципы проектирования

Принцип Реализация
Иммутабельность Все Pydantic-модели frozen — нет утечек мутабельного состояния
Типобезопасность Полные type hints, строгая проверка mypy
Fail-Fast Ошибки валидации ловятся рано с понятными сообщениями
Расширяемость Плагин-архитектура для custom-валидаторов
Производительность Оптимизация для пакетной обработки, минимум аллокаций
Безопасность Sandbox для custom-выражений

Структура модулей

quality-validator/
├── src/kruma_validator/
│   ├── __init__.py          # Публичные экспорты
│   ├── validator.py         # QualityValidator
│   ├── rules.py             # RuleEngine и Rule
│   ├── contract.py          # Загрузка контрактов
│   ├── models.py            # ValidationResult, DLQRecord
│   ├── kafka_consumer.py    # Kafka-интеграция
│   └── examples/
│       └── basic_usage.py
├── tests/
│   ├── test_validator.py
│   ├── test_rules.py
│   ├── test_models.py
│   └── test_contract.py
└── docs/

Ответственности компонентов

Компонент Ответственность
ContractLoader Загрузка YAML-файлов, разрешение внешних ссылок, парсинг метаданных и правил
Contract Иммутабельное представление контракта: topic prefix, имена топиков (raw/prod/dlq)
QualityValidator Высокоуровневый API — валидация одной записи и пакета
RuleEngine Выполнение правил, поддержание batch state (uniqueness), dispatch по типу правила
Rule Pydantic-модель конфигурации правила с параметрами по типам
ValidationResult Результат валидации одной записи: is_valid, errors, warnings
BatchValidationResult Агрегированный результат: prod/dlq-записи, статистика, breakdown ошибок
DLQRecord Обёртка для невалидной записи с полным контекстом ошибок
QualityValidatorConsumer Kafka consumer с автоматической маршрутизацией prod/dlq

Потоки данных

Детальный поток валидации

                     Входные записи
┌─────────────────────────────────────────────────────┐
│              QualityValidator                        │
│                                                      │
│   validate_batch()                                   │
│       │                                              │
│       ├── Для каждой записи → RuleEngine             │
│       │       │                                      │
│       │       ├── Правило 1 (NOT_NULL) → ✓ или ✗    │
│       │       ├── Правило 2 (RANGE)    → ✓ или ✗    │
│       │       ├── Правило 3 (REGEX)    → ✓ или ✗    │
│       │       └── Правило N            → ✓ или ✗    │
│       │       │                                      │
│       │       ▼                                      │
│       │   ValidationResult                           │
│       │   ├── is_valid: bool                         │
│       │   ├── errors: [...]                          │
│       │   └── warnings: [...]                        │
│       │                                              │
│       ▼                                              │
│   BatchValidationResult                              │
└──────────┬───────────────────────┬───────────────────┘
           │                       │
           ▼                       ▼
    ┌──────────────┐       ┌──────────────┐
    │ is_valid=True│       │is_valid=False│
    │ → Prod topic │       │ → DLQ topic  │
    └──────────────┘       │ + Error Meta │
                           └──────────────┘

Dispatch правил

RuleEngine._apply_rule()
┌────────────────────┐
│ Получить значение  │
│ поля (dot notation)│
└────────┬───────────┘
    ┌────┼────┬────┬────┬────┐
    ▼    ▼    ▼    ▼    ▼    ▼
NOT_NULL UNIQUE RANGE REGEX CUSTOM ...
    │    │     │     │     │
    └────┴─────┴─────┴─────┘
  ValidationError или None

Диаграммы классов

Основные классы

┌───────────────────────────┐
│    QualityValidator       │
├───────────────────────────┤
│ - contract: Contract      │
│ - source_topic: str       │
│ - rule_engine: RuleEngine │
├───────────────────────────┤
│ + from_contract_path()    │
│ + from_namespace()        │
│ + validate(record)        │
│ + validate_batch(records) │
│ + get_topics()            │
│ + get_owner_contact()     │
└─────────┬─────────────────┘
          │ uses
┌───────────────────────────┐
│       RuleEngine          │
├───────────────────────────┤
│ - rules: list[Rule]       │
│ - _unique_values: dict    │
├───────────────────────────┤
│ + validate_record(record) │
│ + reset_batch()           │
│ - _apply_rule()           │
│ - _validate_not_null()    │
│ - _validate_unique()      │
│ - _validate_range()       │
│ - _validate_regex()       │
│ - _validate_enum()        │
│ - _validate_freshness()   │
│ - _validate_format()      │
│ - _validate_custom()      │
└───────────────────────────┘

Перечисления

┌───────────────────┐    ┌────────────────┐
│     RuleType      │    │    Severity     │
├───────────────────┤    ├────────────────┤
│ NOT_NULL          │    │ ERROR          │
│ UNIQUE            │    │ WARNING        │
│ RANGE             │    │ INFO           │
│ REGEX             │    └────────────────┘
│ ENUM              │
│ FRESHNESS         │
│ FORMAT            │
│ CUSTOM            │
│ SQL               │
│ REFERENCE         │
└───────────────────┘

Sequence-диаграммы

Kafka Consumer Flow

sequenceDiagram
    participant K as Kafka Consumer
    participant QVC as QVConsumer
    participant QV as QualityValidator
    participant P as Kafka Producer

    K->>QVC: poll() → messages
    QVC->>QV: validate_batch(records)
    QV-->>QVC: BatchValidationResult

    loop Валидные записи
        QVC->>P: send(prod_topic, record)
    end

    loop Невалидные записи
        QVC->>P: send(dlq_topic, dlq_record)
    end

    P->>P: flush()
    QVC->>K: commit()

Точки расширения

Custom-валидатор

from kruma_validator.rules import RuleEngine, Rule
from kruma_validator.models import ValidationError

class ExtendedRuleEngine(RuleEngine):
    """RuleEngine с custom-логикой валидации."""

    def __init__(self, rules: list[Rule], reference_data: dict = None):
        super().__init__(rules)
        self.reference_data = reference_data or {}

    def _validate_reference(self, rule, value, record):
        """Валидация ссылочной целостности."""
        if value is None:
            return None
        table = rule.reference_table
        lookup = {row[rule.reference_field] for row in self.reference_data.get(table, [])}
        if value not in lookup:
            return self._create_error(rule, f"value in {table}", value)
        return None

Плагин-архитектура

from abc import ABC, abstractmethod

class ValidatorPlugin(ABC):
    """Базовый класс для плагинов валидации."""

    @property
    @abstractmethod
    def rule_type(self) -> str:
        """Тип правила, который обрабатывает плагин."""
        ...

    @abstractmethod
    def validate(self, rule, value, record) -> ValidationError | None:
        """Валидация значения по правилу."""
        ...

# Пример: ML-аномалии
class MLAnomalyPlugin(ValidatorPlugin):
    rule_type = "ml_anomaly"

    def __init__(self, model_path: str):
        self.model = load_model(model_path)

    def validate(self, rule, value, record):
        score = self.model.predict(record)
        if score > rule.threshold:
            return ValidationError(...)
        return None

Потокобезопасность

Компонент Состояние Thread Safety
RuleEngine _unique_values Не потокобезопасен
QualityValidator Содержит RuleEngine Не потокобезопасен
QualityValidatorConsumer Consumer/Producer Не потокобезопасен

Рекомендация: используйте thread-local экземпляры или пул валидаторов:

import threading

_local = threading.local()

def get_validator():
    if not hasattr(_local, "validator"):
        _local.validator = QualityValidator.from_contract_path(
            "contracts/domains/sales/orders"
        )
    return _local.validator

Производительность

Throughput (Apple M1, single thread)

Сценарий Записей/сек
Простые правила (3–5) 50 000+
Сложные правила (10+) 20 000+
С REGEX-паттернами 15 000+
С CUSTOM-выражениями 10 000+

Оптимальные настройки

  • Порядок правил: быстрые (NOT_NULL, ENUM) перед медленными (REGEX, CUSTOM)
  • Размер батча: 100–1000 записей — оптимальный баланс
  • UNIQUE-правило: сбрасывайте tracking между независимыми батчами
  • CUSTOM-выражения: держите простыми; сложная логика — через custom-валидаторы

Конфигурация окружения

# Путь к контрактам
export CONTRACTS_PATH=/app/contracts

# Kafka
export KAFKA_BOOTSTRAP_SERVERS=kafka:9092
export CONSUMER_GROUP=quality-validator

# Производительность
export BATCH_SIZE=500
export BATCH_TIMEOUT_MS=2000

# Логирование
export LOG_LEVEL=INFO