Потоковая CDC репликации данных из любой БД в Greenplum с помощью RabbitMQ и Debezium
Аналитика данных в режиме реального времени и плавная интеграция данных имеют решающее значение для принятия взвешенных бизнес-решений в современном мире, основанном на данных. Захват изменения данных (CDC) - это фундаментальная технология захвата и распространения изменений данных из различных баз данных в режиме реального времени. В этой статье мы поговорим о том, как обеспечить CDC в режиме реального времени из любой БД в хранилище данных Greenplum с помощью RabbitMQ и Debezium.
В этой статье мы также поговорим и о том, как Greenplum, инновационное и высокопроизводительное хранилище данных, предназначенное для аналитики Big Data и ML, позволяет раскрыть всю мощь больших данных в режиме реального времени. Кроме того, мы рассмотрим его расширенные возможности потоковой передачи данных, которые позволяют получать, обрабатывать и анализировать данные в режиме реального времени.
Мы также рассмотрим техническую конфигурацию и пошаговые инструкции по настройке этого конвейера данных.
Понятие захвата изменения данных (CDC)
CDC - это подход, позволяющий определять, захватывать и передавать измененные, добавленные и удаленные данные в каких-либо источниках данных. Он отслеживает операции вставки, обновления и удаления в таблицах БД и преобразует их в события. Другие системы могут использовать этот поток событий для различных целей, включая аналитику данных в режиме реального времени, хранение данных, а также интеграцию данных.
RabbitMQ
RabbitMQ – это распределенный брокер сообщений с открытым исходным кодом, который обеспечивает эффективную доставку сообщений в рамках сложных сценариев маршрутизации. Этот инструмент называется «распределенным», потому что обычно работает как кластер узлов, где очереди распределяются (реплицируются) по узлам для обеспечения высокой доступности и отказоустойчивости.
Разработчики используют RabbitMQ для обработки высокопроизводительных и надежных фоновых заданий, а также для интеграции и взаимодействия внутри приложений и между ними. Инструмент применяется для выполнения сложной маршрутизации к консьюмерам и интеграции нескольких приложений и служб с нетривиальной логикой маршрутизации.
RabbitMQ идеально подходит для веб-серверов, которым требуется быстрый запрос-ответ. Этот инструмент распределяет нагрузку между рабочими приложениями при высокой нагрузке (более 20 000 сообщений в секунду) и может обрабатывать фоновые задания или длительные задачи, такие как преобразование PDF, сканирование файлов или масштабирование изображений.
Используя RabbitMQ, мы с легкостью можем отделить источник данных от хранилища данных, обеспечивая тем самым лучшую масштабируемость, надежность и отказоустойчивость системы.
Debezium для CDC
Debezium – это представитель категории ПО CDC, а если точнее — это набор коннекторов для различных СУБД, совместимых с фреймворком Apache Kafka Connect.
Это open source-проект, использующий лицензию Apache License v2.0 и спонсируемый компанией Red Hat. Разработка ведётся с 2016 года и на данный момент в нем представлена официальная поддержка следующих СУБД: MySQL, PostgreSQL, MongoDB, SQL Server.
Обычно Debezium используется, чтобы позволить различным приложениям почти немедленно реагировать на изменение данных в СУБД: события вставки, обновления и удаления, включая отправку push-уведомлений на одно или несколько мобильных устройств, агрегацию изменений и генерацию потока исправлений для объектов. Debezium распределяет процессы мониторинга или коннекторы между несколькими узлами, реплицируя события, чтобы минимизировать риск потери информации.
Возможности потоковой передачи данных Greenplum
Greenplum - хранилище данных с массивно-параллельной обработкой (MPP) данных на базе PostgreSQL. Данная система характеризуется расширенными функциями в области организации потоковой передачи данных, которая позволяет получать данные из внешних источников в режиме реального времени.
С помощью Greenplum Streaming Server (GPSS) организации могут эффективно обрабатывать большие объемы данных и эффективно интегрировать их непосредственно в Greenplum.
В системе CDC, работающей в режиме реального времени GPSS играет важную роль в получении данных из RabbitMQ, перехватываемых Debezium, и переносе событий CDC в таблицы Greenplum для операций INSERT, UPDATE или DELETE.
Являясь неотъемлемой частью экосистемы Greenplum, GPSS предоставляет непревзойденные возможности потоковой передачи данных, которые позволяют оптимизировать конвейеры обработки данных и повысить эффективность аналитических процессов. Он способен обрабатывать огромные объемы потоковых данных (при обработке 10 миллионов событий в секунду с более чем 500 миллиардами строк в многотриллионной базе данных Greenplum и при выполнении ML).
CDC из любой БД в Greenplum в режиме реального времени
В качестве исходной базы данных мы будем использовать PostgreSQL и запустим CDC в режиме реального времени; в конфигурации буду задействовать следующие компоненты:
- БД PostgreSQL 15
- Сервер Debezium 2.4
- RabbitMQ 3.12.2
- Greenplum 6.24 ( + GPSS 1.10.1)
Для того чтобы облегчить выполнение этой демо-версии, все компоненты будут развернуты с помощью Docker-compose:
version: "3.9"
services:
gpdb:
image: docker.io/ahmedrachid/gpdb_demo:6.21
depends_on:
- rabbitmq
privileged: true
entrypoint: /usr/lib/systemd/systemd
ports:
- 5433:5432
volumes:
- ${PWD}/greenplum-db-6.24.0-rhel8-x86_64.rpm:/home/gpadmin/greenplum-db-6.24.0-rhel8-x86_64.rpm
- ${PWD}/gpss-gpdb6-1.10.1-rhel8-x86_64.gppkg:/home/gpadmin/gpss-gpdb6-1.10.1-rhel8-x86_64.gppkg
- ${PWD}/script_gpdb.sh:/home/gpadmin/script_gpdb.sh
rabbitmq:
image: rabbitmq:3-management-alpine
container_name: rabbitmq
ports:
- 5672:5672
- 15672:15672
- 5552:5552
environment:
RABBITMQ_DEFAULT_PASS: root
RABBITMQ_DEFAULT_USER: root
RABBITMQ_DEFAULT_VHOST: vhost
postgres:
image: quay.io/debezium/example-postgres:2.1
container_name: postgres
ports:
- 5432:5432
environment:
- POSTGRES_USER=postgres
- POSTGRES_PASSWORD=postgres
debezium-server:
image: quay.io/debezium/server:2.4
container_name: debezium-server
ports:
- 8080:8080
volumes:
- ./conf:/debezium/conf
depends_on:
- rabbitmq
- gpdb
- postgres
Основной конфигурационный файл Debezium – это conf/application.properties:
debezium.sink.type=rabbitmq debezium.sink.rabbitmq.connection.host=rabbitmq debezium.sink.rabbitmq.connection.port=5672 debezium.sink.rabbitmq.connection.username=root debezium.sink.rabbitmq.connection.password=root debezium.sink.rabbitmq.connection.virtual.host=vhost debezium.sink.rabbitmq.connection.port=5672 debezium.source.connector.class=io.debezium.connector.postgresql.PostgresConnector debezium.source.offset.storage.file.filename=data/offsets.dat debezium.source.offset.flush.interval.ms=0 debezium.source.database.hostname=postgres debezium.source.database.port=5432 debezium.source.database.user=postgres debezium.source.database.password=postgres debezium.source.database.dbname=postgres debezium.source.topic.prefix=tutorial debezium.source.table.include.list=inventory.customers debezium.source.plugin.name=pgoutput debezium.source.tombstones.on.delete=false debezium.sink.rabbitmq.routingKey=inventory_customers
Как видите, мы используем Debezium для передачи событий CDC из таблицы PostgreSQL под названием inventory.customers в кластер RabbitMQ.
Развертывание базы данных PostgreSQL
Выполните следующую команду Docker для того, чтобы установить и развернуть предварительно сконфигурированный контейнер PostgreSQL:
docker compose up postgres -d
Теперь Вы можете подключиться к базе данных и изучить таблицу, которую мы используем для потоковой передачи событий CDC:
$ docker-compose exec postgres env PGOPTIONS="--search_path=inventory" bash -c 'psql -U $POSTGRES_USER postgres -c "SELECT * FROM customers;"' id | first_name | last_name | email ------+------------+-----------+----------------------- 1001 | Sally | Thomas | sally.thomas@acme.com 1002 | George | Bailey | gbailey@foobar.com 1003 | Edward | Walker | ed@walker.com 1004 | Anne | Kretchmar | annek@noanswer.org (4 rows)
Установка и настройка RabbitMQ:
Чтобы развернуть кластер RabbitMQ, мы используем следующую команду docker compose:
docker compose up rabbitmq -d
После того как Вы развернули RabbitMQ с помощью Docker Compose, Вы можете настроить его в соответствии с Вашими потребностями.
С помощью веб-интерфейса RabbitMQ Вы можете настроить работу сервера RabbitMQ.
Чтобы получить доступ к управлению, откройте веб-браузер и перейдите на IP-адрес или доменное имя Вашего экземпляра RabbitMQ, а затем на номер порта управления (по умолчанию 15672) - http://localhost:15672.
Теперь Вы можете войти в систему, используя учетные данные по умолчанию (root/root) или имя пользователя и пароль, указанные в файле Docker Compose.
В системе управления необходимо создать следующее:
- Новый обменник RabbitMQ Exchange под названием tutorial.inventory.customers;
- Новое потоковое расширение RabbitMQ Stream под названием inventory.customers, привязанное к обменнику tutorial.inventory.customers помощью routingKey inventory_customers
Установка и настройка сервера Debezium
Теперь, когда наши контейнеры PostgreSQL и RabbitMQ запущены, пришло время запустить сервер Debezium. Для этого выполним следующие действия:
docker compose up debezium-server -d
После этого Вы должны увидеть новое открытое соединение с кластером RabbitMQ:
Вы также должны увидеть новые сообщения в потоке RabbitMQ:
Настройка Greenplum для принятия событий CDC в режиме реального времени
Чтобы обеспечить CDC из RabbitMQ в Greenplum в режиме реального времени мы можем использовать Greenplum 6.24, предварительно настроенный с помощью Greenplum Streaming Server 1.10.1.
Во-первых, нужно развернуть контейнер Greenplum:
docker compose up gpdb -d
Затем для загрузки событий CDC нужно создать таблицу customers:
$ docker compose exec -ti gpdb bash $ su - gpadmin $ psql postgres -c 'CREATE TABLE public.customers (id INT, first_name TEXT, last_name TEXT, email TEXT) DISTRIBUTED BY (id);'
GPSS или Greenplum Streaming Server эффективно обрабатывает потоки данных, поступающие из Kafka и RabbitMQ в базу данных Greenplum.
Гибкая, масштабируемая архитектура обеспечивает высокопроизводительный прием данных с минимальными задержками. GPSS разработан для работы с различными форматами данных, включая TEXT, CSV, JSON, Avro и другие, что делает его подходящим для CDC в режиме реального времени.
Чтобы запустить загрузку событий CDC из потока RabbitMQ, необходимо запустить процесс GPSS:
nohup gpss &
Затем создайте задание GPSS , используя конфигурацию, приведенную ниже:
DATABASE: postgres USER: gpadmin HOST: localhost PORT: 5432 VERSION: 2 RABBITMQ: INPUT: SOURCE: SERVER: root:root@rabbitmq:5552 STREAM: inventory.customers VIRTUALHOST: vhost DATA: COLUMNS: - NAME: j TYPE: json FORMAT: json ERROR_LIMIT: 25 OUTPUT: TABLE: customers MODE: MERGE MATCH_COLUMNS: - id DELETE_CONDITION: ((j->>'payload')::json->>'op')='d' MAPPING: - NAME: id EXPRESSION: CASE WHEN ((j->>'payload')::json->>'op')='d' THEN (((j->>'payload')::json->>'before')::json->>'id')::int ELSE (((j->>'payload')::json->>'after')::json->>'id')::int END - NAME: first_name EXPRESSION: CASE WHEN ((j->>'payload')::json->>'op')='d' THEN (((j->>'payload')::json->>'before')::json->>'first_name')::text ELSE (((j->>'payload')::json->>'after')::json->>'first_name')::text END - NAME: last_name EXPRESSION: CASE WHEN ((j->>'payload')::json->>'op')='d' THEN (((j->>'payload')::json->>'before')::json->>'last_name')::text ELSE (((j->>'payload')::json->>'after')::json->>'last_name')::text END - NAME: email EXPRESSION: CASE WHEN ((j->>'payload')::json->>'op')='d' THEN (((j->>'payload')::json->>'before')::json->>'email')::text ELSE (((j->>'payload')::json->>'after')::json->>'email')::text END
Отправьте конфигурацию задания GPSS:
$ gpsscli submit gpss_job.yaml
Запустите задание GPSS, чтобы начать обработку данных из потока RabbitMQ:
$ gpsscli start gpss_job
После запуска Вы должны увидеть, что выполнение задания инициировано:
[gpadmin@df7ac82e5633 ~]$ gpsscli list JobName JobID GPHost GPPort DataBase Schema Table Topic Status gpss_job 7f40703e01734ad570c31b028fdfc615 localhost 5432 postgres public customers JOB_RUNNING
Добавьте данные в таблицу customers (PostgreSQL)
Теперь вы должны увидеть данные (4 записи), поступающие в таблицу Greenplum. Вы можете сгенерировать несколько событий CDC при помощи операций INSERT/UPDATE или DELETE в исходной базе данных.
Добавляем данные:
INSERT INTO customers (id, first_name, last_name, email) SELECT i, 'First_' || i, 'Last_' || i, 'first' || i || '.last' || i || '@example.com' AS email FROM generate_series(1010, 1000004) AS i;
Данные CDC, обрабатываемые Greenplum в режиме реального времени
- После выполнения команды INSERT INTO события CDC оперативно применяются к таблице customers (Greenplum) в режиме реального времени.
- С другой стороны, Вы можете обновить запись в исходной базе данных PostgreSQL:
В результате Вы можете увидеть, что событие UPDATE было успешно применено и на стороне Greenplum, что означает то, что данные были согласованы и синхронизированы.
Заключение
Возможности потоковой передачи данных Greenplum, RabbitMQ Streams и Debezium в совокупности образуют надежный и эффективный CDC-конвейер обработки данных, работающий в режиме реального времени.
Такая система позволяет организациям анализировать данные в режиме реального времени и принимать взвешенные решения, основанные на данных. Используя все возможности потоковой передачи данных Greenplum, компании могут получить весомое конкурентное преимущество и обогнать своих соперников в гонке за производительность и прибылью.














