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

Интеграция с PostgreSQL

⚠️ Критически важный принцип

Приложение СНАЧАЛА сохраняет данные в PostgreSQL, ПОТОМ отправляет в Data Platform!

┌─────────────────┐     ┌─────────────────┐     ┌─────────────────┐     ┌─────────────────┐
│  Application    │────▶│  PostgreSQL     │────▶│  ETL Script     │────▶│  API Gateway    │
│                 │     │  (COMMIT!)      │     │  (Avro + mTLS)  │     │  (mTLS)         │
└─────────────────┘     └─────────────────┘     └─────────────────┘     └─────────────────┘
                              ↑                                                   │
                       СНАЧАЛА COMMIT!                                           ▼
                                                                        ┌─────────────────┐
                                                                        │  Kafka *.raw    │
                                                                        └─────────────────┘

Почему именно так?

  1. Надёжность: Данные не потеряются даже если Data Platform недоступна
  2. Аудит: Все данные сохранены в вашей БД
  3. Переотправка: Можно переотправить данные при сбое
  4. 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 сертификат:

  1. Создайте тикет в JIRA с запросом на сертификат
  2. Укажите контракт: warehouse/inventory
  3. Укажите владельца (должен совпадать с owner в контракте)
  4. Data Platform team выдаст сертификат
  5. Сохраните сертификат и ключ в безопасное место

См. подробнее: docs/implementation/10_mtls_certificates.md

Рекомендации

  1. Используйте batch ETL для большинства случаев (простота + надёжность)
  2. Используйте NOTIFY/LISTEN для event-driven сценариев
  3. Используйте CDC только если нужен real-time с минимальной нагрузкой на БД
  4. Всегда добавляйте индексы на exported_to_platform
  5. Мониторьте количество неотправленных записей

Мониторинг

-- Сколько записей ожидает отправки?
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/