Интеграция с PostgreSQL¶
⚠️ Критически важный принцип¶
Приложение СНАЧАЛА сохраняет данные в PostgreSQL, ПОТОМ отправляет в Data Platform!
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ Application │────▶│ PostgreSQL │────▶│ ETL Script │────▶│ API Gateway │
│ │ │ (COMMIT!) │ │ (Avro + mTLS) │ │ (mTLS) │
└─────────────────┘ └─────────────────┘ └─────────────────┘ └─────────────────┘
↑ │
СНАЧАЛА COMMIT! ▼
┌─────────────────┐
│ Kafka *.raw │
└─────────────────┘
Почему именно так?¶
- Надёжность: Данные не потеряются даже если Data Platform недоступна
- Аудит: Все данные сохранены в вашей БД
- Переотправка: Можно переотправить данные при сбое
- Ownership: Вы владеете данными, Data Platform — это дополнительная копия
Обзор¶
PostgreSQL — наиболее распространённая СУБД. Этот пример показывает интеграцию PostgreSQL с системой контрактов данных через Apache Avro и mTLS.
Варианты интеграции¶
1. Batch ETL (простой и надёжный)¶
Периодическая выгрузка данных батчами.
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ PostgreSQL │────▶│ Python ETL │────▶│ API Gateway │
│ (SELECT *) │ │ (Avro + mTLS) │ │ (mTLS) │
└─────────────────┘ └─────────────────┘ └─────────────────┘
Преимущества: - Простота настройки - Контроль над трансформацией - Автоматическая ретрансляция при сбое - Подходит для batch workloads
Запуск:
export DB_HOST=app-postgres.local
export DB_PASSWORD=secret
export MTLS_CERT=/secure/certs/warehouse.inventory.crt
export MTLS_KEY=/secure/certs/warehouse.inventory.key
python batch_etl.py \
--table inventory_updates \
--contract warehouse/inventory \
--version 1.0.0
Запуск по расписанию (cron):
# Каждые 5 минут
*/5 * * * * /usr/bin/python3 /app/batch_etl.py --table inventory_updates --contract warehouse/inventory
2. NOTIFY/LISTEN (real-time)¶
Real-time реакция на изменения через PostgreSQL NOTIFY.
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ PostgreSQL │────▶│ Python Listener│────▶│ API Gateway │
│ (NOTIFY) │ │ (Avro + mTLS) │ │ (mTLS) │
└─────────────────┘ └─────────────────┘ └─────────────────┘
Преимущества: - Встроенная функциональность PostgreSQL - Низкая латентность (< 1 секунды) - Не требует дополнительной инфраструктуры - Event-driven архитектура
Создание триггера:
CREATE OR REPLACE FUNCTION notify_inventory_update()
RETURNS TRIGGER AS $$
BEGIN
PERFORM pg_notify(
'inventory_updates',
json_build_object(
'update_id', NEW.update_id,
'sku', NEW.sku
)::text
);
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER inventory_update_trigger
AFTER INSERT ON inventory_updates
FOR EACH ROW
EXECUTE FUNCTION notify_inventory_update();
Запуск слушателя:
export DB_HOST=app-postgres.local
export DB_PASSWORD=secret
export MTLS_CERT=/secure/certs/warehouse.inventory.crt
export MTLS_KEY=/secure/certs/warehouse.inventory.key
python notify_listener.py \
--channel inventory_updates \
--table inventory_updates \
--contract warehouse/inventory \
--version 1.0.0
Запуск как systemd service:
[Unit]
Description=PostgreSQL NOTIFY/LISTEN → Data Platform
After=postgresql.service
[Service]
Type=simple
User=dataplatform
Environment="DB_HOST=localhost"
Environment="MTLS_CERT=/secure/certs/warehouse.inventory.crt"
Environment="MTLS_KEY=/secure/certs/warehouse.inventory.key"
ExecStart=/usr/bin/python3 /app/notify_listener.py \
--channel inventory_updates \
--table inventory_updates \
--contract warehouse/inventory
Restart=always
RestartSec=10
[Install]
WantedBy=multi-user.target
3. CDC через Debezium (для OLTP с высокими требованиями)¶
Если нужен real-time с минимальной нагрузкой на БД.
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ PostgreSQL │────▶│ Debezium │────▶│ API Gateway │
│ (WAL logs) │ │ + Avro Converter│ │ (mTLS) │
└─────────────────┘ └─────────────────┘ └─────────────────┘
Преимущества: - Минимальная нагрузка на БД (читает WAL) - Захват всех изменений (INSERT, UPDATE, DELETE) - Real-time с низкой latency
Настройка PostgreSQL:
-- Включить logical replication
ALTER SYSTEM SET wal_level = 'logical';
ALTER SYSTEM SET max_replication_slots = 4;
ALTER SYSTEM SET max_wal_senders = 4;
-- Создать пользователя для Debezium
CREATE USER debezium WITH REPLICATION LOGIN PASSWORD 'secure_password';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO debezium;
-- Создать publication
CREATE PUBLICATION orders_publication FOR TABLE
orders,
order_items;
Подготовка таблицы¶
Добавьте флаг для отслеживания экспорта:
ALTER TABLE inventory_updates
ADD COLUMN IF NOT EXISTS exported_to_platform BOOLEAN DEFAULT false,
ADD COLUMN IF NOT EXISTS exported_at TIMESTAMP WITH TIME ZONE;
CREATE INDEX idx_exported ON inventory_updates (exported_to_platform);
Получение mTLS сертификата¶
Для работы с API Gateway нужен mTLS сертификат:
- Создайте тикет в JIRA с запросом на сертификат
- Укажите контракт:
warehouse/inventory - Укажите владельца (должен совпадать с owner в контракте)
- Data Platform team выдаст сертификат
- Сохраните сертификат и ключ в безопасное место
См. подробнее: docs/implementation/10_mtls_certificates.md
Рекомендации¶
- Используйте batch ETL для большинства случаев (простота + надёжность)
- Используйте NOTIFY/LISTEN для event-driven сценариев
- Используйте CDC только если нужен real-time с минимальной нагрузкой на БД
- Всегда добавляйте индексы на
exported_to_platform - Мониторьте количество неотправленных записей
Мониторинг¶
-- Сколько записей ожидает отправки?
SELECT COUNT(*)
FROM inventory_updates
WHERE exported_to_platform = false;
-- Самая старая неотправленная запись
SELECT MIN(created_at)
FROM inventory_updates
WHERE exported_to_platform = false;
Troubleshooting¶
Ошибка: Certificate not found¶
# Проверьте наличие сертификата
ls -la /secure/certs/warehouse.inventory.crt
# Проверьте permissions (должен быть 600)
chmod 600 /secure/certs/warehouse.inventory.key
Ошибка: mTLS authentication failed¶
Убедитесь что: 1. CN в сертификате совпадает с топиком (warehouse.inventory) 2. Сертификат не истёк 3. CA файл актуальный
# Проверить CN в сертификате
openssl x509 -in /secure/certs/warehouse.inventory.crt -noout -subject
# Проверить срок действия
openssl x509 -in /secure/certs/warehouse.inventory.crt -noout -enddate
Ошибка: Database connection failed¶
# Проверьте подключение к БД
psql -h $DB_HOST -U $DB_USER -d $DB_NAME -c "SELECT 1"
# Проверьте переменные окружения
env | grep DB_
Ещё примеры¶
- 1C integration:
../1c/ - Mindbox webhooks:
../mindbox/