Использование плагина PostgreSQL pgoutput для сбора данных об изменениях с помощью Debezium в Azure
В этой статье мы вкратце расскажем о том, как работает плагин pgoutput. Я не буду повторять множество деталей и использую контейнерные версии (с помощью Docker Compose) для Kafka connect, Kafka (и Zookeeper), чтобы все было максимально просто. Итак, единственное, что Вам нужно, это Azure PostgreSQL, который Вы можете настроить с помощью различных опций, включая Azure Portal, Azure CLI, Azure PowerShell, ARM.
Ресурсы доступны на GitHub — https://github.com/abhirockzz/debezium-postgres-pgoutput
Использование правильного publication.autocreate.mode
При использовании плагина pgoutput важно, чтобы Вы использовали соответствующее значение для publication.autocreate.mode. Если Вы используете all_tables (по умолчанию), Вам нужно убедиться, что публикация создана заранее для конкретной таблицы (таблиц), которую Вы хотите настроить для захвата данных об изменениях. Если публикация не найдена, коннектор попытается создать ее с помощью CREATE PUBLICATION <имя_публикации> ДЛЯ ВСЕХ ТАБЛИЦ, что приведет к неудаче из-за отсутствия разрешений.
Два других варианта работают, как и ожидалось:
- disabled: необходимо убедиться, что публикация создана заранее. Коннектор не будет пытаться создать публикацию, если при запуске не обнаружится, что она существует - он выбросит исключение и остановится.
- filtered: вы можете (по желанию) выбрать создание публикации заранее. Если публикация не будет найдена, коннектор создаст новую публикацию для всех таблиц, соответствующих текущей конфигурации фильтра.
Все это было описано здесь:
https://debezium.io/documentation/reference/1.3/connectors/postgresql.html#postgresql-on-azure
Предлагаю попробовать разные сценарии
Исходный вариант:
git clone https://github.com/abhirockzz/debezium-postgres-pgoutput && cd debezium-postgres-pgoutput
Запустите контейнеры Kafka, Zookeeper и Kafka Connect:
export DEBEZIUM_VERSION=1.2 docker-compose up
Первое включение контейнеров может занять некоторое время.
Когда все контейнеры будут запущены, подключитесь к Azure PostgreSQL, создайте таблицу и вставьте в нее данные следующим образом:
psql -h <DBNAME>.postgres.database.azure.com -p 5432 -U <DBUSER>@<DBNAME> -W -d postgres --set=sslmode=requirepsql -h abhishgu-pg.postgres.database.azure.com -p 5432 -U abhishgu@abhishgu-pg -W -d postgres --set=sslmode=requireCREATE TABLE inventory (id SERIAL, item VARCHAR(30), qty INT, PRIMARY KEY(id));
Если для параметра publication.autocreate.mode установлено значение filtered
Это хорошо работает с Azure PostgreSQL - для этого не требуются права суперпользователя, поскольку коннектор создает публикацию для определенной таблицы (таблиц) на основе значений фильтра/*списка.
Обновите файл конфигурации коннектора (pg-source-connector.json), указав в нем сведения о Вашем экземпляре Azure PostgreSQL, а затем создайте коннектор.
Чтобы создать коннектор:
curl -X POST -H "Content-Type: application/json" --data @pg-source-connector.json http://localhost:8083/connectors
Обратите внимание на журналы (в терминале docker compose):
Creating new publication 'mytestpub' for plugin 'PGOUTPUT' [io.debezium.connector.postgresql.connection.PostgresReplicationConnection] Once the connector starts, check the publications in PostgreSQL: pubname | schemaname | tablename -----------+------------+----------- ytestpub | public | inventory
Работает?
Вставьте несколько записей в таблицу inventory:
psql -h <DBNAME>.postgres.database.azure.com -p 5432 -U <DBUSER>@<DBNAME> -W -d postgres --set=sslmode=requireINSERT INTO inventory (item, qty) VALUES ('apples', '100');
INSERT INTO inventory (item, qty) VALUES ('oranges', '42');select * from inventory;
Коннектор должен передавать события изменений из PostgreSQL WAL (журнал опережающей записи) в Kafka. Проверьте сообщения в соответствующей теме Kafka:
//exec into the kafka docker container docker exec -it debezium-postgres-pgoutput_kafka_1 bashcd bin && ./kafka-console-consumer.sh --topic myserver.public.inventory --bootstrap-server kafka:9092 --from-beginnin
Вы должны увидеть несколько полезных нагрузок событий журнала изменений (соответствующих двум INSERT).
Да, они достаточно подробные, поскольку схема включена в полезную нагрузку.
Измените publication.autocreate.mode на disabled
Для этого режима нам нужна публикация, созданная заранее. Поскольку у нас уже есть одна (mytestpub), просто используйте ее. Все, что Вам нужно сделать, это обновить publication.autocreate.mode в pg-source-connector.json на disabled.
Создайте коннектор заново:
//delete curl -X DELETE localhost:8083/connectors/inventory-connector//create curl -X POST -H "Content-Type: application/json" --data @pg-source-connector.json http://localhost:8083/connectors
Протестируйте его от конца до конца, выполнив те же шаги, что и в предыдущем разделе, - все должно работать отлично!
Для подтверждения обновите publication.name в конфигурации коннектора на несуществующее. Коннектор не запустится из-за отсутствия публикации (как и ожидалось)
Попробуйте publication.autocreate.mode = all_tables
Установите publication.autocreate.mode на all_tables, publication.name на несуществующее (например, testpub1) и создайте коннектор:
curl -X POST -H "Content-Type: application/json" --data @pg-source-connector.json http://localhost:8083/connectors
(как и ожидалось) Он завершится с ошибкой, подобной этой:
....
INFO Creating new publication 'testpub1' for plugin 'PGOUTPUT' (io.debezium.connector.postgresql.connection.PostgresReplicationConnection:127) ERROR WorkerSourceTask{id=inventory-connector-0} Task threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runtime.WorkerTask:179) io.debezium.jdbc.JdbcConnectionException: ERROR: must be superuser to create FOR ALL TABLES publication ....
Обратите внимание, что для создания публикации FOR ALL TABLES необходимо быть суперпользователем - как уже говорилось ранее, CREATE PUBLICATION <publication_name> FOR ALL TABLES; не удалось из-за отсутствия прав суперпользователя.
Как я уже говорил, Вам нужно обойти эту проблему, создав публикацию вручную только для определенных таблиц.
Очистка
Для очистки удалите экземпляр Azure PostgreSQL с помощьюaz postgres server delete и удалите контейнеры:
az postgres server delete -g <resource group> -n <server name>docker-compose down -v




