
Перед командой, выросшей на 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:
-
снимает консистентный снепшот текущих данных (
snapshot.mode: initial); -
дальше читает бинлоги MariaDB построчно и стримит каждое изменение (insert/update/delete) в Kafka;
-
на другом конце 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 |
|
|
Debezium Kafka |
|
|
Debezium Connect |
|
Ниже — пример для переноса db-source-01 (MariaDB) → db-target-01 (PostgreSQL).
Шаг 1. Инфраструктура: ZooKeeper, Kafka, Kafka Connect
docker network create debezium-netdocker run -d --name zookeeper \ --network debezium-net \ -p 2181:2181 -p 2888:2888 -p 3888:3888 \ quay.io/debezium/zookeeper:3.1.1.Finaldocker 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.Finaldocker 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.Finaldocker 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/bashtables=( '{"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 5done
На практике коннектор чаще всего падает в 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.ymlsource_db_host: source-db.example.internalsink_db_host: sink-db.example.internaltables: - 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.
ссылка на оригинал статьи https://habr.com/ru/articles/1061438/