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

Руководство по интеграции

Установка, настройка и развёртывание Kruma Quality Validator — от локального использования до production-деплоя в Kubernetes.

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


Установка

Из PyPI

pip install kruma-validator

Из исходников

git clone http://192.168.88.72/kruma/quality-validator
cd quality-validator
pip install -e ".[dev]"

Зависимости

Обязательные:

  • Python >= 3.11
  • pydantic >= 2.0
  • pyyaml >= 6.0

Опциональные (Kafka-интеграция):

  • kafka-python-ng >= 2.2.0
  • avro >= 1.11.0

Проверка установки

from kruma_validator import QualityValidator, __version__
print(f"kruma-validator: {__version__}")
# kruma-validator: 1.0.1

Быстрый старт

from kruma_validator import QualityValidator

# 1. Создать валидатор из контракта
validator = QualityValidator.from_contract_path(
    contract_path="contracts/domains/sales/orders",
    contracts_base_path="/path/to/contracts"
)

# 2. Валидировать запись
result = validator.validate(record)

if result.is_valid:
    print("Запись валидна")
else:
    for error in result.errors:
        print(f"Ошибка: {error.rule_name}{error.message}")

# 3. Пакетная валидация
batch_result = validator.validate_batch(records)
prod_records = batch_result.get_prod_records()   # → .prod
dlq_records = batch_result.get_dlq_records()     # → .dlq

Альтернативные способы инициализации

# По namespace и name
validator = QualityValidator.from_namespace(
    namespace="sales", name="orders",
    contracts_base_path="/path/to/contracts"
)

# Из объекта Contract
from kruma_validator import ContractLoader
loader = ContractLoader("/path/to/contracts")
contract = loader.load("domains/sales/orders")
validator = QualityValidator(contract)

Именование топиков

BREAKING CHANGE в v1.0.1

С версии 1.0.1 все топики используют точку как разделитель:

Было (v1.0.0) Стало (v1.0.1)
sales.orders_prod sales.orders.prod
sales.orders_dlq sales.orders.dlq
sales.orders.raw sales.orders.raw (без изменений)

Формат

{namespace}.{name}.{suffix}
Суффикс Назначение Пример
.raw Сырые входящие данные sales.orders.raw
.prod Провалидированные данные sales.orders.prod
.dlq Dead Letter Queue sales.orders.dlq

Программное получение

topics = validator.get_topics()
# {"raw": "sales.orders.raw", "prod": "sales.orders.prod", "dlq": "sales.orders.dlq"}

Kafka Consumer

Базовый consumer

from kruma_validator.kafka_consumer import QualityValidatorConsumer

consumer = QualityValidatorConsumer(
    contract_path="contracts/domains/sales/orders",
    bootstrap_servers="kafka-1:9092,kafka-2:9092,kafka-3:9092",
    consumer_group="quality-validator-sales-orders",
    batch_size=100,
    batch_timeout_ms=1000,
    contracts_base_path="/app/contracts"
)

# Блокирующий запуск
consumer.run()

Consumer с обработчиками

def handle_batch(batch_result):
    contract = f"{batch_result.contract_namespace}.{batch_result.contract_name}"
    print(f"{contract}: {batch_result.valid_count}/{batch_result.total_count} валидных")

consumer.run(on_batch=handle_batch)

Расширенная конфигурация Kafka

consumer = QualityValidatorConsumer(
    contract_path="contracts/domains/sales/orders",
    bootstrap_servers="kafka:9092",
    consumer_group="quality-validator",
    batch_size=500,
    batch_timeout_ms=2000,
    consumer_config={
        "security_protocol": "SASL_SSL",
        "sasl_mechanism": "SCRAM-SHA-256",
        "sasl_plain_username": "validator",
        "sasl_plain_password": "secret",
    },
    producer_config={
        "acks": "all",
        "retries": 3,
        "compression_type": "lz4",
    }
)

Docker

Запуск

docker build -t kruma/quality-validator:latest .

docker run -d \
  --name quality-validator-orders \
  -e KAFKA_BOOTSTRAP_SERVERS=kafka-1:9092,kafka-2:9092 \
  -e CONTRACT_PATH=domains/sales/orders \
  -e CONSUMER_GROUP=quality-validator-sales-orders \
  -v /path/to/contracts:/contracts:ro \
  kruma/quality-validator:latest

Docker Compose

services:
  quality-validator-orders:
    image: kruma/quality-validator:latest
    environment:
      KAFKA_BOOTSTRAP_SERVERS: kafka:9092
      CONTRACT_PATH: domains/sales/orders
      CONSUMER_GROUP: qv-sales-orders
      CONTRACTS_PATH: /contracts
      LOG_LEVEL: INFO
    volumes:
      - ./contracts:/contracts:ro
    depends_on:
      - kafka
    restart: unless-stopped
    deploy:
      resources:
        limits:
          memory: 512M
          cpus: "0.5"

  quality-validator-inventory:
    image: kruma/quality-validator:latest
    environment:
      KAFKA_BOOTSTRAP_SERVERS: kafka:9092
      CONTRACT_PATH: domains/warehouse/inventory
      CONSUMER_GROUP: qv-warehouse-inventory
      CONTRACTS_PATH: /contracts
    volumes:
      - ./contracts:/contracts:ro
    restart: unless-stopped

Kubernetes

apiVersion: apps/v1
kind: Deployment
metadata:
  name: quality-validator-orders
spec:
  replicas: 2
  selector:
    matchLabels:
      app: quality-validator-orders
  template:
    spec:
      containers:
        - name: validator
          image: kruma/quality-validator:latest
          env:
            - name: KAFKA_BOOTSTRAP_SERVERS
              valueFrom:
                configMapKeyRef:
                  name: kafka-config
                  key: bootstrap-servers
            - name: CONTRACT_PATH
              value: domains/sales/orders
          resources:
            requests: { memory: "256Mi", cpu: "250m" }
            limits: { memory: "512Mi", cpu: "500m" }

Переменные окружения

Переменная Описание По умолчанию
KAFKA_BOOTSTRAP_SERVERS Адреса Kafka-брокеров localhost:9092
CONTRACT_PATH Путь к директории контракта
CONTRACTS_PATH Базовый путь для всех контрактов .
CONSUMER_GROUP ID consumer group quality-validator
BATCH_SIZE Записей в батче 100
BATCH_TIMEOUT_MS Макс. ожидание батча (мс) 1000
LOG_LEVEL Уровень логирования INFO
METRICS_PORT Порт Prometheus-метрик 9090

Prometheus-метрики

Настройка

from prometheus_client import Counter, Histogram, start_http_server

RECORDS_PROCESSED = Counter(
    "validator_records_processed_total",
    "Total records processed",
    ["contract", "status"]
)

VALIDATION_ERRORS = Counter(
    "validator_validation_errors_total",
    "Validation errors by rule",
    ["contract", "rule_name"]
)

# Запуск сервера метрик
start_http_server(9090)

Рекомендуемые метрики

Метрика Тип Описание
validator_records_processed_total Counter Обработано записей по статусу
validator_validation_errors_total Counter Ошибки по правилам
validator_validation_duration_seconds Histogram Время валидации батча
validator_batch_size Gauge Текущий размер батча
validator_dlq_records_total Counter Записей в DLQ

Grafana-запросы

# Pass rate
sum(rate(validator_records_processed_total{status="valid"}[5m])) by (contract)
/ sum(rate(validator_records_processed_total[5m])) by (contract)

# Топ ошибающихся правил
topk(10, sum(rate(validator_validation_errors_total[1h])) by (rule_name))

# Throughput
sum(rate(validator_records_processed_total[5m])) by (contract)

CI/CD интеграция

GitLab CI

validate-contracts:
  stage: validate
  image: python:3.12-slim
  script:
    - pip install kruma-validator
    - python -m kruma_validator.validate_contracts --path contracts/
  rules:
    - changes:
        - contracts/**/*

test-validator:
  stage: test
  image: python:3.12-slim
  script:
    - pip install -e ".[dev]"
    - pytest tests/ -v --cov=kruma_validator --cov-report=xml
  coverage: '/TOTAL.*\s+(\d+%)/'

Pre-commit hook

# .pre-commit-config.yaml
repos:
  - repo: local
    hooks:
      - id: validate-contracts
        name: Validate Data Contracts
        entry: python -c "from kruma_validator import ContractLoader; ContractLoader().load('contracts/domains/sales/orders')"
        language: python
        files: ^contracts/
        additional_dependencies: [kruma-validator]

Тестирование

Unit-тесты

import pytest
from kruma_validator import QualityValidator

@pytest.fixture
def validator():
    return QualityValidator.from_contract_path("contracts/domains/sales/orders")

def test_valid_record(validator, valid_record):
    result = validator.validate(valid_record)
    assert result.is_valid is True

def test_missing_order_id(validator, valid_record):
    del valid_record["order_id"]
    result = validator.validate(valid_record)
    assert result.is_valid is False
    assert any(e.rule_name == "order_id_required" for e in result.errors)

Property-based тесты

from hypothesis import given, strategies as st
from kruma_validator.rules import RuleEngine, Rule, RuleType

@given(st.text(min_size=1))
def test_not_null_passes_on_non_none(value):
    rule = Rule(name="test", field="f", rule_type=RuleType.NOT_NULL, severity="error")
    engine = RuleEngine([rule])
    errors = engine.validate_record({"f": value})
    assert len(errors) == 0

Устранение неполадок

Проблема Причина Решение
ImportError: kafka-python Нет Kafka-зависимостей pip install kafka-python-ng
FileNotFoundError: Contract Неверный путь Проверить CONTRACTS_PATH и CONTRACT_PATH
Consumer не обрабатывает Неверная consumer group Проверить CONSUMER_GROUP и имена топиков
Высокий DLQ rate Строгие правила Пересмотреть severity (error → warning)

Debug-режим

export LOG_LEVEL=DEBUG