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 на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс по Apache Airflow и NiFi » Настройка пайплайна с использованием Airflow и PostgreSQL

Настройка пайплайна с использованием Airflow и PostgreSQL

В этой статье мы развернем ETL пайплайн с использованием Apache Airflow, RabbitMQ и PostgreSQL на локальной машине.

 

Шаг 1: Установка необходимых инструментов

Перед началом убедитесь, что на вашем компьютере установлены следующие инструменты:

  • Docker Desktop — для управления контейнерами.
  • PyCharm — для разработки Python-кода.
  • PgAdmin4 — для удобного управления PostgreSQL.

 

После установки этих инструментов создайте файл docker-compose.yml, в котором будут прописаны все образы, необходимые для развертывания нашего пайплайна. Вы можете скачать готовый файл с настройками или создать свой собственный.

 

Шаг 2: Развертывание сервисов с помощью Docker

Создайте новый проект в PyCharm, загрузите в корневую папку docker-compose.yml файл и выполните следующую команду в терминале:

docker compose up -d

Эта команда автоматически скачает необходимые образы и запустит контейнеры для Airflow, RabbitMQ, PostgreSQL, Redis и MongoDB. Как только контейнеры будут запущены, Docker настроит необходимую структуру папок в вашем проекте.

 

Шаг 3: Настройка соединений между Airflow и PostgreSQL

Теперь подключимся к PostgreSQL через PgAdmin. Введите логин и пароль: airflow и airflow. В PgAdmin у вас должны появиться две базы данных, одна из которых будет использоваться для работы с Airflow.

 

Создайте таблицу, которая будет использоваться для тестирования ETL пайплайна.

Чтобы настроить соединение между Airflow и PostgreSQL, откройте веб-интерфейс Airflow через Docker Desktop. Войдите с логином и паролем. В разделе Connections создайте новое соединение с типом Postgres, указав хост вашего контейнера с PostgreSQL и данные для подключения.

 

Шаг 4: Создание простого DAG для вставки данных в PostgreSQL

Для проверки соединения с базой данных создадим простой DAG в Airflow, который будет вставлять данные в таблицу PostgreSQL. Для этого установите необходимые зависимости через PyCharm:

pip install apache-airflow apache-airflow-providers-postgres

 

Пример кода DAG для вставки данных:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.postgres.hooks.postgres import PostgresHook

hook = PostgresHook(postgres_conn_id='Postgres')
request = "insert into public.data (name, birth_date, age) values ('Anton', '2004-11-14', 19)"

def insert():
    hook.run(request)

with DAG(dag_id="base_dag", default_args={"owner": "Kirill"}) as dag:
    t = PythonOperator(task_id='insert_postgres', python_callable=insert)

 

Когда этот код будет готов, в интерфейсе Airflow появится новый DAG. Запустите его, и вы увидите, что данные были успешно добавлены в вашу таблицу PostgreSQL.

 

Шаг 5: Настройка соединения между Airflow и RabbitMQ

Теперь настраиваем RabbitMQ. Перейдите в интерфейс RabbitMQ, доступный через Docker Desktop, и войдите с логином и паролем guest и guest. В разделе Queues создайте очередь с названием queue.airflow.

В PyCharm установите библиотеку airflow-provider-rabbitmq для работы с RabbitMQ:

pip install apache-airflow-providers-rabbitmq

 

Теперь добавьте соединение с RabbitMQ в Airflow, указав имя контейнера RabbitMQ как rabbitmq_main.

Шаг 6: Пример работы с RabbitMQ в Airflow

Создадим задачу, которая будет получать сообщения из RabbitMQ и передавать их в другие таски. В папке dags создадим папку tasks_broker и добавим два файла: sensor.py и get_data.py.

sensor.py:

from rabbitmq_provider.sensors.rabbitmq import RabbitMQSensor

sensor = RabbitMQSensor(
    task_id="sensor",
    queue_name="queue.airflow",
    rabbitmq_conn_id="RabbitMQ",
)

 

get_data.py:

import json
from airflow.decorators import task

@task(task_id="get_data")
def get_data(**kwargs) -> None:
    message = json.loads(kwargs['task_instance'].xcom_pull(task_ids='sensor'))
    print("#########################################################################################")
    print(message)
    print("#########################################################################################")

 

Создадим новый DAG для взаимодействия с RabbitMQ:

from airflow import DAG
from tasks_broker.sensor import sensor
from tasks_broker.get_data import get_data
with DAG(dag_id="python_broker", schedule_interval=None) as dag:
    sensor >> get_data()

 

Шаг 7: Реализация ETL процесса

Для реализации ETL процесса создадим:

  1. Извлечение данных из RabbitMQ — получаем данные в формате JSON.
  2. Обработка данных — выбираем нужные поля из полученных данных.
  3. Загрузка данных в PostgreSQL — сохраняем обработанные данные в базе.

 

Создадим файлы: extract.py, transformation.py, и load.py.

extract.py (получение данных):

from rabbitmq_provider.sensors.rabbitmq import RabbitMQSensor

extract = RabbitMQSensor(
    task_id="extract",
    queue_name="queue.airflow",
    rabbitmq_conn_id="RabbitMQ",
)

 

transformation.py (обработка данных):

import json
from airflow.decorators import task

@task(task_id="transformation")
def transformation(**kwargs) -> dict:
    message = json.loads(kwargs['task_instance'].xcom_pull(task_ids='extract'))
    new_message = {"name": message["name"], "birth_date": message["birth_date"], "age": message["age"]}
    return new_message

 

load.py (загрузка в PostgreSQL):

from airflow.decorators import task
from airflow.providers.postgres.hooks.postgres import PostgresHook

@task(task_id="load")
def load(**kwargs) -> None:
    message = kwargs['task_instance'].xcom_pull(task_ids='transformation')
    hook = PostgresHook(postgres_conn_id='Postgres')
    request = f"insert into public.data (name, birth_date, age) values ('{message['name']}', '{message['birth_date']}', {message['age']})"
    hook.run(request)

 

Теперь создадим основной файл python_broker_postgres.py:

from airflow import DAG
from tasks_broker_postgres.extract import extract
from tasks_broker_postgres.transformation import transformation
from tasks_broker_postgres.load import load

with DAG(dag_id="python_broker_postgres", schedule_interval=None) as dag:
    extract >> transformation() >> load()

 

Запустив DAG, вы сможете увидеть, как данные из RabbitMQ обрабатываются и загружаются в PostgreSQL.

В этой статье мы научились настраивать и развертывать ETL пайплайн с использованием Apache Airflow, RabbitMQ и PostgreSQL. Мы создали DAG для обработки данных и их загрузки в базу данных, а также настроили необходимые соединения между компонентами. Это решение позволит автоматизировать и ускорить обработку данных в различных проектах.

 

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

← Предыдущая статья
Как мы организуем 2000+ моделей DBT в Apache Airflow
Следующая статья →
Ассеты в оркестраторах данных: концепция и реализация в Airflow

Решения

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

Клиенты
  • ООО "Уральская транспортная компания" — это транспортно-логистическая компания, специализирующаяся на железнодорожных перевозках грузов, создана в 2009 году.

  • Авиакомпания NordStar (АО «АК «НордСтар») – работает под данным брендом с 2008 г. и сейчас входит в топ-15 крупнейших российских авиакомпаний (данные Росавиации) с пассажирооборотом более 1 млн человек в год. АО «АК «НордСтар» выполняет и внутренние, и внешние рейсы, а ее основные хабы - Домодедово, Пулково и Емельяново. С 2021 года компания является базовым перевозчиком аэропорта Норильск.

  • ООО "Интернэшнл Ресторант Брэндс" – это крупнейший франчайзинговый партнер компании Yum! Brands Russia & CIS в России, отвечающий за рост и развитие бренда KFC на территории РФ. На сегодняшний день у компании более 350 ресторанов. Ежедневно в рестораны приходит 200 000+ гостей.

  • Розничный и интернет-магазин 12 Storeez один из лидеров на рынке женской одежды. С географией рынка не только на территории России, своя продукция представлена еще и в таких странах как Казахстан и Дубай.

  • Решения
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • 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 и политикой конфиденциальности.