Как интегрировать Greenplum с Apache Airflow
Доступ и обработка данных Greenplum в Apache Airflow с помощью драйвера CData JDBC.
Apache Airflow поддерживает создание, планирование и мониторинг рабочих процессов инженерии данных. В паре с драйвером CData JDBC Driver for Greenplum Airflow может работать с живыми данными Greenplum. В этой статье описывается то, как подключиться к данным Greenplum из экземпляра Apache Airflow, запросить их и сохранить результаты в CSV-файле.
Благодаря встроенной оптимизированной обработке данных драйвер CData JDBC обеспечивает непревзойденную производительность при взаимодействии с живыми данными Greenplum. Когда вы отправляете сложные SQL-запросы к Greenplum, драйвер передает поддерживаемые SQL-операции, такие как фильтры и агрегации, непосредственно в Greenplum и использует встроенный SQL-движок для обработки неподдерживаемых операций на стороне клиента (часто это SQL-функции и операции JOIN). Встроенные динамические запросы к метаданным позволяют работать с данными Greenplum и анализировать их, используя собственные типы данных.
Настройка подключения к Greenplum
Встроенный конструктор строк подключения
Для помощи в создании URL-адреса JDBC используйте конструктор строк подключения, встроенный в драйвер Greenplum JDBC Driver. Либо дважды щелкните JAR-файл, либо запустите jar-файл из командной строки.
java -jar cdata.jdbc.greenplum.jar
Заполните свойства соединения и скопируйте строку соединения в буфер обмена.
Чтобы подключиться к Greenplum, задайте свойства соединения Server, Port (порт по умолчанию - 5432) и Database, а также задайте User и Password, которые вы хотите использовать для аутентификации на сервере. Если свойство Database не указано, будет использоваться база данных по умолчанию для аутентифицируемого пользователя.
Чтобы разместить драйвер JDBC в кластерной среде или в облаке, Вам потребуется лицензия (полная или пробная), а также ключ времени выполнения (RTK).
Ниже перечислены основные свойства, необходимые для установления JDBC-соединения:
|
Свойство |
Значение |
|
URL соединение с БД |
jdbc:greenplum:RTK=5246...;User=user;Password=admin;Database=dbname;Server=127.0.0.1;Port=5432; |
|
Название класса доайвера БД |
cdata.jdbc.greenplum.GreenplumDriver |
Установление JDBC-соединения в Airflow
- Войдите в свой экземпляр Apache Airflow.
- На панели навигации Вашего экземпляра Airflow наведите курсор на Admin, а затем нажмите Connections
- Далее нажмите знак + на следующем экране, чтобы создать новое соединение.
- В форме добавления соединения заполните необходимые свойства соединения:
- Connection id: Имя соединения, например: greenplum_jdbc
- Connection Type: JDBC-соединение
- Connection URL: URL JDBC-соединения, указанный выше, например: jdbc:greenplum:RTK=5246...;User=user;Password=admin;Database=dbname;Server=127.0.0.1;Port=5432;)
- Driver Class: cdata.jdbc.greenplum.GreenplumDriver
- Driver Path: PATH/TO/cdata.jdbc.greenplum.jar
- Проверьте новое соединение, нажав кнопку Test в нижней части формы.
- После сохранения нового соединения на новом экране должен появиться зеленый баннер, сообщающий, что в список соединений была добавлена новая строка:
Создание DAG
Группа DAG в Airflow - это сущность, которая хранит процессы, необходимые для рабочего процесса, и может быть запущена для его выполнения. Наш рабочий процесс состоит в том, чтобы просто выполнить SQL-запрос к данным Greenplum максимально оперативно и сохранить результаты в CSV-файле.
- В домашнем каталоге должна быть папка «airflow». Внутри нее мы можем создать новую директорию и назвать ее «dags». Здесь мы будем хранить файлы Python, которые преобразуются в DAGы Airflow, отображаемые в пользовательском интерфейсе.
- Далее создайте новый Python-файл и назовите его greenplum_hook.py. Вставьте в этот новый файл следующий код:
import time
from datetime import datetime
from airflow.decorators import dag, task
from airflow.providers.jdbc.hooks.jdbc import JdbcHook
import pandas as pd
# Declare Dag
@dag(dag_id="greenplum_hook", schedule_interval="0 10 * * *", start_date=datetime(2022,2,15), catchup=False, tags=['load_csv'])
# Define Dag Function
def extract_and_load():
# Define tasks
@task()
def jdbc_extract():
try:
hook = JdbcHook(jdbc_conn_id="jdbc")
sql = """ select * from Account """
df = hook.get_pandas_df(sql)
df.to_csv("/{some_file_path}/{name_of_csv}.csv",header=False, index=False, quoting=1)
# print(df.head())
print(df)
tbl_dict = df.to_dict('dict')
return tbl_dict
except Exception as e:
print("Data extract error: " + str(e))
jdbc_extract()
sf_extract_and_load = extract_and_load()
3. Сохраните этот файл и обновите экземпляр Airflow. В списке групп DAG Вы должны увидеть новую группу DAG под названием «greenplum_hook».
4. Щелкните по этой группе DAG, на новом экране нажмите на переключатель unpause, чтобы он стал синим, а затем нажмите на кнопку запуска (т. е. воспроизведения), чтобы запустить группу DAG. Это выполнит SQL-запрос в нашем файле greenplum_hook.py и экспортирует результаты в формате CSV в файл, путь к которому мы указали в нашем коде.
5. После запуска нашей новой группы DAG мы проверяем папку Downloads (или место, которое Вы выбрали в своем Python-сценарии) и видим, что CSV-файл создан - в данном случае account.csv.
6.Откройте файл CSV и убедитесь в том, что Ваши данные Greenplum теперь доступны для использования в формате CSV благодаря Apache Airflow.












