Руководство по интеграции¶
Установка, настройка и развёртывание Kruma Quality Validator — от локального использования до production-деплоя в Kubernetes.
Версия: 1.0.1 | Владелец: Kruma Platform Team
Установка¶
Из PyPI¶
Из исходников¶
Зависимости¶
Обязательные:
- 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 (без изменений) |
Формат¶
| Суффикс | Назначение | Пример |
|---|---|---|
.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) |