Настройка пайплайна с использованием 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 процесса создадим:
- Извлечение данных из RabbitMQ — получаем данные в формате JSON.
- Обработка данных — выбираем нужные поля из полученных данных.
- Загрузка данных в 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 для обработки данных и их загрузки в базу данных, а также настроили необходимые соединения между компонентами. Это решение позволит автоматизировать и ускорить обработку данных в различных проектах.





