BI Consult Desktop Logo BI Consult Mobile Logo
  • Russian BI Исследование российских bi
  • Перейти на Fine BI
  • Контакты
  • +7 812 334-08-01
    +7 499 608-13-06
  • Отправить сообщение
  • Главная
  • Продукты Эксперт-BI
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Сельское хозяйство
    • Энергетика
    • FMCG
    • Девелоперы
    • Маркетплейсы
    • Пищевая промышленность
    • Фармацевтика
    • Построение Data Platform
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и FP&A
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • IBP
    • ИТ (CIO)
    • Закупки
  • Платформы
    • Системы бизнес-анализа (BI)
    • Интегрированное бизнес-планирование (IBP)
    • Хранилища данных (DWH / Lakehouse)
    • Каталоги данных (Data Catalog)
    • Системы ETL и ELT
    • AI / Исскуственный интеллект
    • Шина данных (ESB)
    • Система управления мастер-данными (MDM)
    • Семантический слой
  • Услуги
    • Переход на отечественные BI и DWH системы
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений и DWH
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Курсы
    • Учебный курс Информационная грамотность (Data Literacy)
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Greenplum
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt (Data Build Tool)
  • Компания
    • Руководство
    • Новости
    • Клиенты
    • Карьера
    • Скачать
    • Контакты

BI

  • FineBI
  • FineReport
  • FineDataLink
  • FineChatBI (FineAI)
  • Коннекторы данных из 1С в BI
  • Airflow / Nifi
  • Visiology
  • PIX BI
  • Modus BI
  • Yandex.DataLens
  • Open-source BI: Superset/Metabase
  • Luxms BI
  • AW BI + Alpha BI
  • FlyBI + Форсайт. Аналитическая Платформа
  • Loginom
  • Триафлай
  • AI / Исскуственный интеллект
  • Optimacros
  • Навигатор BI
  • Семантический слой

СУБД

  • Arenadata
  • ClickHouse
  • Greenplum
  • Postgres Professional
  • TData

Другое

  • Построение Data Platform
    • Аналитическое хранилище данных
    • Data Lake и Data Engineering
    • Подробнее про Data Lake
    • Внедрение Lakehouse
      • Apache Doris
      • StarRocks
      • Trino
    • Миграция витрин из пропиетарных DWH на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс по Greenplum » Потоковая CDC репликации данных из любой БД в Greenplum с помощью RabbitMQ и Debezium

Потоковая 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, компании могут получить весомое конкурентное преимущество и обогнать своих соперников в гонке за производительность и прибылью.

 

Узнать стоимость решенияЗапросить видео презентацию

← Предыдущая статья
Greenplum от Vmware – это векторная база, отлично подходящая для аналитики Ваших данных
Следующая статья →
От MySQL к Greenplum: уроки и выводы, сделанные в процессе миграции от одной БД к другой

Решения

Анализировать ФинансыУвеличивайте ПродажиОптимальный Склад и ЛогистикаМаркетинговые Метрики

Клиенты
  • Компания ООО "Комус" - один из лидеров российского рынка оптовых продаж офисных товаров и техники. Компания поставляет широкий ассортимент продукции - от канцелярских принадлежностей до компьютерной техники и офисной мебели.

  • ЭГИС - международная фармацевтическая компания, основанная в 1907 году в Венгрии. Компания имеет представительства более чем в 60 странах мира, в том числе в России. Компания ЭГИС является одним из ведущих производителей дженерических лекарственных средств в Центральной и Восточной Европе. Её деятельность охватывает все звенья производственно-сбытовой фармацевтической цепочки.

  • КАМИ – компания-лидер по поставкам тяжёлых станков в России, занимающаяся продажей и обслуживанием оборудования для обработки металла и дерева, изготовления мебели и не только. На сегодняшний день в компании работают более 1300 человек, запущено 10 обучающих центров, в продаже более 7000 единиц техники. 

  • ПАО «Транснефть» – крупнейшая российская нефтепроводная компания. «Транснефть» обеспечивает транспортировку более 85% добываемых в России нефти и нефтепродуктов.

  • Решения
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Энергетика
    • Фармацевтика
  • Услуги
    • Переход на отечественные BI и DWH
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Техническая поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Платформы
    • FineBI
    • FineReport
    • FineDataLink
    • Коннекторы данных из 1С в BI
    • Airflow + NiFi
    • Visiology
    • Luxms BI
    • Modus BI
    • PIX BI
    • Arenadata
    • ClickHouse
    • Greenplum
    • Postgres Professional
    • Open-source BI: Superset/Metabase
    • Loginom
    • Yandex.DataLens
    • AI / Исскуственный интеллект
    • Optimacros
    • Шины данных
  • Курсы
    • Учебный курс Информационная грамотность
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt
  • Функциональные решения
    • Создание Data Lake
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и прогнозная аналитика
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • Сквозная аналитика
  • Компания
    • О нас
    • Руководство
    • Новости
    • Клиенты
    • Скачать
    • Контакты
    • Политика конфиденциальности
RutubeVkontakteLinkedInYouTube
ООО "Би Ай Консалт",
ИНН: 7811437757,
ОГРН: 1097847154184
199178, Россия,
Санкт-Петербург,
6-ая линия В.О., Д. 63, 4 этаж
Тел: +7 (812) 334-08-01
Тел: +7 (499) 608-13-06
E-mail: info@biconsult.ru

 

 

 

 

 

×

Пользуясь сайтом, вы соглашаетесь с использованием cookies и политикой конфиденциальности.