Перед командой, выросшей на MySQL/MariaDB, встаёт вопрос перехода на PostgreSQL: лицензии, экосистема расширений, более предсказуемое поведение при конкурентных нагрузках - причин достаточно. Проблема в другом: как перенести боевую базу, которая не останавливается ни на минуту, без часов простоя и без риска потерять данные, записанные во время переноса.

В этой статье - как мы решали эту задачу для интернет-магазина на MariaDB, почему готовые консольные конвертеры не годятся для «живой» миграции, и как выглядит рабочая схема на Debezium + Kafka Connect, включая Ansible-роль для повторяемого запуска.

Все имена хостов, баз, топиков и учётные данные в примерах - вымышленные.

Почему не подошли готовые утилиты

Первая мысль при слове «миграция MySQL → PostgreSQL» - взять один из известных конвертеров и прогнать через него дамп. Мы попробовали несколько вариантов, и у всех обнаружился один и тот же фундаментальный недостаток: это утилиты одноразового переноса, а не репликации. Они снимают снепшот на момент запуска и не умеют донакатывать изменения, случившиеся в источнике после начала работы. Для базы, которая не останавливается, это означает окно даунтайма на время переноса - для нас неприемлемое.

Плюс к этому у каждой утилиты нашлись собственные болячки:

  • pgloader - самый популярный вариант, но в его issue-трекере регулярно всплывают падения по памяти (heap exhaustion) на больших таблицах, ошибки парсинга (ESRAP-PARSE-ERROR), проблемы с «нулевыми» датами MySQL, конфликты имён при превышении лимита PostgreSQL в 63 символа и дублирующимися именами индексов, которые MySQL допускает неявно, а PostgreSQL - нет (пример, ещё один, и ещё). На части наших таблиц миграция просто зависала на середине.

  • pg_chameleon - ближе к тому, что нам было нужно (реальная репликация через чтение бинлогов), но требует binlog_format=ROW, обязательного primary key на каждой таблице, а при ошибке загрузки строки просто выбрасывает конфликтную таблицу из репликации - то есть часть данных молча перестаёт синхронизироваться, и это легко пропустить. Кроме того, направление «PostgreSQL → MySQL» у него экспериментальное и сильно ограниченное - жизнеспособна только миграция в одну сторону.

  • py-mysql2pgsql - по сути, заброшенный проект: релизов нет уже несколько лет, поддержка неактивна, для рабочей нагрузки не рассматривали.

Ни один из этих инструментов не даёт того, что было нужно: непрерывной синхронизации источника и приёмника, чтобы можно было мигрировать данные заранее, дать таблицам «дореплицироваться» и в момент отключения приложения от MariaDB переключить его на PostgreSQL буквально с разницей в секунды.

Решение: Debezium как CDC-платформа

Debezium - это набор коннекторов для Kafka Connect, реализующих Change Data Capture (CDC). В отличие от разовых конвертеров, Debezium:

  1. снимает консистентный снепшот текущих данных (snapshot.mode: initial);

  2. дальше читает бинлоги MariaDB построчно и стримит каждое изменение (insert/update/delete) в Kafka;

  3. на другом конце sink-коннектор разбирает поток из Kafka и применяет изменения в PostgreSQL через upsert.

В результате PostgreSQL-реплика непрерывно «догоняет» источник, и переключение приложения можно делать в любой удобный момент, когда лаг репликации близок к нулю - без остановки MariaDB на время переноса.

Плата за это - инфраструктурная сложность (нужен Kafka с ZooKeeper, Kafka Connect, диск под очередь) и то, что все изменения временно материализуются в Kafka в виде JSON - в нашем тесте перенос данных занял около 50 ГБ дискового пространства именно из-за этого формата. Разворачивать стек лучше на отдельной машине, а не на сервере с боевой базой.

Версии компонентов

Компонент

Версия

MariaDB

v11.7.2

PostgreSQL

v15.12 (Debian 15.12-0+deb12u2)

Debezium ZooKeeper

quay.io/debezium/zookeeper:3.1.1.Final

Debezium Kafka

quay.io/debezium/kafka:3.1.1.Final

Debezium Connect

quay.io/debezium/connect:3.1.1.Final

Ниже - пример для переноса db-source-01 (MariaDB) → db-target-01 (PostgreSQL).

Шаг 1. Инфраструктура: ZooKeeper, Kafka, Kafka Connect

docker network create debezium-net

docker run -d --name zookeeper \
  --network debezium-net \
  -p 2181:2181 -p 2888:2888 -p 3888:3888 \
  quay.io/debezium/zookeeper:3.1.1.Final

docker run -d --name kafka \
  --network debezium-net \
  -p 9092:9092 \
  -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
  -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092 \
  -e ZOOKEEPER_CONNECT=zookeeper:2181 \
  -e KAFKA_AUTO_CREATE_TOPICS_ENABLE=true \
  quay.io/debezium/kafka:3.1.1.Final

docker run -d --name connect \
  --network debezium-net \
  -p 8083:8083 \
  -e BOOTSTRAP_SERVERS=kafka:9092 \
  -e GROUP_ID=1 \
  -e CONFIG_STORAGE_TOPIC=my_connect_configs \
  -e OFFSET_STORAGE_TOPIC=my_connect_offsets \
  -e STATUS_STORAGE_TOPIC=my_connect_statuses \
  quay.io/debezium/connect:3.1.1.Final

docker exec -it kafka /kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka:9092 --create \
  --topic schemachanges-example \
  --partitions 1 --replication-factor 1

Отдельный топик schemachanges-example нужен Debezium для хранения истории изменений схемы источника - без него source-коннектор не запустится.

Шаг 2. Source-коннектор (MariaDB)

curl -i -X POST -H "Content-Type:application/json" http://localhost:8083/connectors/ -d '{
  "name": "mariadb-connector",
  "config": {
    "connector.class": "io.debezium.connector.mariadb.MariaDbConnector",
    "database.hostname": "db-source-01.example.int",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "<MARIADB_PASSWORD>",
    "database.server.id": "1",
    "database.include.list": "example_shop",
    "database.connectionTimeZone": "Europe/Moscow",
    "topic.prefix": "db-source-01-example-int",
    "schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
    "schema.history.internal.kafka.topic": "schemachanges-example",
    "include.schema.changes": "true",
    "snapshot.mode": "initial",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.delete.handling.mode": "rewrite",
    "transforms.unwrap.drop.tombstones": "false",
    "max.batch.size": "100",
    "max.queue.size": "500"
  }
}'

curl -s http://localhost:8083/connectors/mariadb-connector/status | jq

Статус (connector.state и tasks[].state) должен быть RUNNING.

Шаг 3. Sink-коннектор (PostgreSQL)

curl -i -X POST -H "Content-Type:application/json" http://localhost:8083/connectors/ -d '{
  "name": "postgres-sink-connector",
  "config": {
    "connector.class": "io.debezium.connector.jdbc.JdbcSinkConnector",
    "tasks.max": "1",
    "topics.regex": "db-source-01-example-int\\.example_shop\\.(?!(device_fingerprints|viewed_products|migration_versions|region_zone|store_cell)$).*",
    "connection.url": "jdbc:postgresql://db-target-01:5432/example_shop?sslmode=disable",
    "connection.username": "postgres",
    "connection.password": "<POSTGRES_PASSWORD>",
    "insert.mode": "upsert",
    "primary.key.mode": "record_key",
    "primary.key.fields": "id",
    "auto.create": "true",
    "auto.evolve": "true",
    "delete.enabled": "true",
    "quote.identifiers": "true",
    "schema.evolution": "basic",
    "transforms": "unwrap,route",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
    "transforms.route.regex": "db-source-01-example-int\\.example_shop\\.(.*)",
    "transforms.route.replacement": "$1"
  }
}'

curl -s http://localhost:8083/connectors/postgres-sink-connector/status | jq

Проблема таблиц без ключа id

Debezium ожидает, что первичный ключ записи и есть ключ Kafka-сообщения, по умолчанию - колонка id. Но в реальных схемах почти всегда находится десяток таблиц с составным или нестандартным уникальным ключом: связки many-to-many, справочники по коду/номеру и т.п. Для них общий sink-коннектор из regex-фильтра выше исключён явно - иначе Debezium попытается писать по несуществующему id и завалит upsert.

В нашем случае таких таблиц набралось порядка 15: например, email_blocklist (ключ - email), store_zone (store_id, zone_id), regions (number) и подобные им по структуре связки и справочники.

Важно. Такую таблицу нужно добавить не только в персональный коннектор, но и в exclude-список общего (тот самый (?!(...)$) в topics.regex из шага 3) - иначе её подхватят оба коннектора и общий упадёт в FAILED.

Для каждой такой таблицы поднимается отдельный sink-коннектор со своим primary.key.fields:

#!/bin/bash

tables=(
  '{"table": "email_blocklist", "keys": "email"}'
  '{"table": "store_zone", "keys": "store_id,zone_id"}'
  '{"table": "regions", "keys": "number"}'
  # ... остальные таблицы с нестандартным ключом
)

for t in "${tables[@]}"; do
  table=$(echo "$t" | jq -r '.table')
  keys=$(echo "$t" | jq -r '.keys')
  connector_name="postgres-sink-${table}-connector"
  topics_regex="db-source-01-example-int\\.example_shop\\.${table}$"

  config=$(cat <<EOF
{
  "name": "${connector_name}",
  "config": {
    "connector.class": "io.debezium.connector.jdbc.JdbcSinkConnector",
    "tasks.max": "1",
    "topics.regex": "${topics_regex}",
    "connection.url": "jdbc:postgresql://db-target-01:5432/example_shop?sslmode=disable",
    "connection.username": "postgres",
    "connection.password": "<POSTGRES_PASSWORD>",
    "insert.mode": "upsert",
    "primary.key.mode": "record_key",
    "primary.key.fields": "${keys}",
    "auto.create": "true",
    "auto.evolve": "true",
    "delete.enabled": "true",
    "quote.identifiers": "true",
    "schema.evolution": "basic",
    "transforms": "unwrap,route",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
    "transforms.route.regex": "db-source-01-example-int\\.example_shop\\.(.*)",
    "transforms.route.replacement": "\$1"
  }
}
EOF
)

  curl -i -X POST -H "Content-Type: application/json" \
    http://localhost:8083/connectors/ -d "${config}"
done

Мониторинг статуса коннекторов

Простейший способ убедиться, что ничего не «упало» и не ушло в FAILED, - опрашивать /status в цикле по всем коннекторам (общему и per-table) и проверять состояние коннектора и каждой его таски:

while true; do
  error_found=false
  for connector in "${all_connectors[@]}"; do
    result=$(curl -s "http://localhost:8083/connectors/${connector}/status")
    [ -z "$result" ] && { echo "ERROR: нет ответа от ${connector}"; error_found=true; continue; }

    state=$(echo "$result" | jq -r '.connector.state')
    [ "$state" != "RUNNING" ] && { echo "ERROR: ${connector} -> ${state}"; error_found=true; }

    for task_state in $(echo "$result" | jq -r '.tasks[].state'); do
      [ "$task_state" != "RUNNING" ] && { echo "ERROR: task ${connector} -> ${task_state}"; error_found=true; }
    done
  done
  sleep 5
done

На практике коннектор чаще всего падает в FAILED из-за несовместимости типов (например, zerodate в MySQL) или из-за потери соединения с БД - оба случая видно сразу по этому циклу, без необходимости лезть в логи Connect.

Продакшен-запуск: Ansible-роль

Ручные curl-запросы удобны для отладки, но неудобны для боевого прогона: легко забыть шаг, опечататься в regex или не заметить, что коннектор ушёл в FAILED. Поэтому весь процесс мы обернули в Ansible-роль - это даёт три вещи, которых нет у набора bash-скриптов:

  • Одна команда на весь цикл. ansible-playbook migrate.yml --tags debezium_up поднимает сеть, контейнеры, топик и все коннекторы (общий + по каждой таблице с нестандартным ключом) за один прогон, без ручного повторения curl по списку таблиц.

  • Идемпотентность и повторный запуск. Роль безопасно перезапускать: uri-модуль сам обрабатывает уже существующие коннекторы (status_code: [201, 409]), а не падает при повторном создании.

  • Переносимость между окружениями. Хосты, пароли, список таблиц и их ключи, исключения - всё вынесено в переменные defaults/main.yml. Под новую пару баз роль адаптируется правкой одного файла, а не кода.

Отдельные теги закрывают весь жизненный цикл: debezium_up (поднять и настроить), debezium_monitor (опросить статус всех коннекторов и вывести сводку) и debezium_down (удалить коннекторы, контейнеры и сеть после завершения миграции).

Пример того, как выглядят переменные конкретной миграции:

# roles/database/migration/defaults/main.yml
source_db_host: source-db.example.internal
sink_db_host: sink-db.example.internal

tables:
  - table: "orders"
    keys: "id"
  - table: "regions"
    keys: "number"
  # ...

Сама роль (таски, шаблоны запросов, цикл мониторинга) - это уже вопрос оформления под конкретный Ansible-проект команды, здесь принципиальна именно идея: обернуть последовательность curl-вызовов в декларативную, идемпотентную и параметризуемую структуру, которую можно запускать и переиспользовать одной командой.

Здесь же удобно закрыть и проблему из предыдущего раздела: excluded_tables для общего коннектора имеет смысл не поддерживать вручную отдельным списком, а собирать из tables (например, через map(attribute='table') в Jinja-шаблоне regex). Тогда список таблиц с нестандартным ключом и exclude-список общего коннектора физически не смогут разъехаться.

На что обратить внимание при переносе на себя

  • Диск. Kafka хранит все изменения в виде JSON, объём быстро растёт - закладывайте disk-запас существенно больше размера самой базы. У нас на тестовом переносе ушло около 50 ГБ при не самой большой базе.

  • Отдельная машина под Debezium-стек. Не разворачивайте Kafka/Connect на сервере с боевой MariaDB - снепшот и так создаёт дополнительную нагрузку на источник.

  • Таблицы без id. Прогоните схему заранее и явно составьте список таблиц с составным/нестандартным ключом - иначе общий sink-коннектор либо не создаст их в PostgreSQL, либо будет писать некорректно.

  • Часовой пояс. database.connectionTimeZone в source-коннекторе стоит явно выставлять под часовой пояс сервера MariaDB - иначе временные поля разъедутся при переносе.

  • schema.evolution: basic. Достаточно для добавления новых колонок «на лету», но не для сложных миграций схемы (переименования, смена типов) - их лучше катить руками до переключения.

  • Переключение приложения. Реальное отключение от MariaDB и переход на PostgreSQL стоит делать только когда лаг репликации (разница между последним событием в бинлоге и последним применённым в Postgres) близок к нулю - иначе часть данных, записанных «в последнюю секунду», рискует потеряться.

Итог

Готовые конвертеры вроде pgloader или pg_chameleon хорошо подходят для разового переноса статичного дампа, но плохо - для миграции живой, постоянно пишущей базы: либо нет догоняющей репликации, либо инструмент молча выбрасывает проблемные таблицы из синхронизации. Связка Debezium + Kafka Connect закрывает именно эту задачу: снепшот плюс непрерывный поток изменений из бинлогов, что даёт возможность мигрировать данные заранее и переключить приложение на новую базу с минимальным (секунды, а не часы) окном рассинхронизации - за счёт более сложной инфраструктуры и заметного расхода диска на промежуточное хранение в Kafka.