Первоначальная реализация
This commit is contained in:
@@ -0,0 +1,100 @@
|
||||
"""
|
||||
DAG для демонстрации условного выполнения задач в Airflow
|
||||
Уровень: Продвинутый
|
||||
"""
|
||||
from datetime import datetime, timedelta
|
||||
from airflow import DAG
|
||||
from airflow.operators.python import PythonOperator, BranchPythonOperator
|
||||
from airflow.operators.dummy import DummyOperator
|
||||
import random
|
||||
|
||||
# Определение DAG
|
||||
default_args = {
|
||||
'owner': 'student',
|
||||
'depends_on_past': False,
|
||||
'start_date': datetime(2023, 1, 1),
|
||||
'email_on_failure': False,
|
||||
'email_on_retry': False,
|
||||
'retries': 1,
|
||||
'retry_delay': timedelta(minutes=5)
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
'branching_dag',
|
||||
default_args=default_args,
|
||||
description='DAG для изучения условного выполнения задач в Airflow',
|
||||
schedule_interval=timedelta(days=1),
|
||||
catchup=False,
|
||||
tags=['educational', 'branching', 'advanced']
|
||||
)
|
||||
|
||||
def check_data_quality():
|
||||
"""Проверка качества данных - случайным образом определяет, какие данные использовать"""
|
||||
# В реальном сценарии здесь будет проверка качества данных
|
||||
# Для учебных целей просто случайное решение
|
||||
quality_score = random.random() # случайное число от 0 до 1
|
||||
|
||||
if quality_score > 0.5:
|
||||
print(f"Качество данных хорошее (оценка: {quality_score:.2f}), используем CSV")
|
||||
return 'process_csv_branch'
|
||||
else:
|
||||
print(f"Качество данных требует внимания (оценка: {quality_score:.2f}), используем JSON")
|
||||
return 'process_json_branch'
|
||||
|
||||
def process_csv_data():
|
||||
"""Обработка CSV данных"""
|
||||
print("Обработка CSV файла...")
|
||||
# Здесь будет логика обработки CSV файла
|
||||
return "CSV данные обработаны"
|
||||
|
||||
def process_json_data():
|
||||
"""Обработка JSON данных"""
|
||||
print("Обработка JSON файла...")
|
||||
# Здесь будет логика обработки JSON файла
|
||||
return "JSON данные обработаны"
|
||||
|
||||
def merge_results():
|
||||
"""Объединение результатов из разных веток"""
|
||||
print("Объединение результатов из разных веток...")
|
||||
return "Результаты объединены"
|
||||
|
||||
# Определение задач
|
||||
start_task = DummyOperator(
|
||||
task_id='start_task',
|
||||
dag=dag
|
||||
)
|
||||
|
||||
check_quality_task = BranchPythonOperator(
|
||||
task_id='check_data_quality',
|
||||
python_callable=check_data_quality,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
process_csv_task = PythonOperator(
|
||||
task_id='process_csv_branch',
|
||||
python_callable=process_csv_data,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
process_json_task = PythonOperator(
|
||||
task_id='process_json_branch',
|
||||
python_callable=process_json_data,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
merge_task = PythonOperator(
|
||||
task_id='merge_results',
|
||||
python_callable=merge_results,
|
||||
trigger_rule='none_failed_or_skipped', # Выполняется, когда одна из веток завершена
|
||||
dag=dag
|
||||
)
|
||||
|
||||
end_task = DummyOperator(
|
||||
task_id='end_task',
|
||||
dag=dag
|
||||
)
|
||||
|
||||
# Установка зависимостей
|
||||
start_task >> check_quality_task
|
||||
check_quality_task >> [process_csv_task, process_json_task]
|
||||
[process_csv_task, process_json_task] >> merge_task >> end_task
|
||||
@@ -0,0 +1,186 @@
|
||||
"""
|
||||
DAG для демонстрации ETL процессов в Airflow
|
||||
Уровень: Средний-Продвинутый
|
||||
"""
|
||||
from datetime import datetime, timedelta
|
||||
from airflow import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.providers.postgres.operators.postgres import PostgresOperator
|
||||
import pandas as pd
|
||||
|
||||
# Определение DAG
|
||||
default_args = {
|
||||
'owner': 'student',
|
||||
'depends_on_past': False,
|
||||
'start_date': datetime(2023, 1, 1),
|
||||
'email_on_failure': False,
|
||||
'email_on_retry': False,
|
||||
'retries': 1,
|
||||
'retry_delay': timedelta(minutes=5)
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
'data_processing_dag',
|
||||
default_args=default_args,
|
||||
description='DAG для изучения ETL процессов в Airflow',
|
||||
schedule_interval=timedelta(days=1),
|
||||
catchup=False,
|
||||
tags=['educational', 'etl', 'intermediate']
|
||||
)
|
||||
|
||||
def create_sample_data():
|
||||
"""Создание примерных данных для ETL процесса"""
|
||||
# Создаем файлы с данными клиентов и заказов
|
||||
customers_data = {
|
||||
'customer_id': [1, 2, 3, 4, 5],
|
||||
'name': ['Alice Johnson', 'Bob Smith', 'Charlie Brown', 'Diana Prince', 'Eve Wilson'],
|
||||
'email': ['alice@example.com', 'bob@example.com', 'charlie@example.com', 'diana@example.com', 'eve@example.com'],
|
||||
'join_date': ['2023-01-15', '2023-02-20', '2023-03-10', '2023-04-05', '2023-05-12']
|
||||
}
|
||||
|
||||
orders_data = {
|
||||
'order_id': [101, 102, 103, 104, 105, 106],
|
||||
'customer_id': [1, 2, 1, 3, 4, 5],
|
||||
'product': ['Laptop', 'Monitor', 'Keyboard', 'Mouse', 'Tablet', 'Headphones'],
|
||||
'amount': [1200, 300, 80, 25, 400, 100],
|
||||
'order_date': ['2023-10-01', '2023-10-02', '2023-10-03', '2023-10-04', '2023-10-05', '2023-10-06']
|
||||
}
|
||||
|
||||
# Сохраняем в CSV
|
||||
pd.DataFrame(customers_data).to_csv('/opt/airflow/data/input/customers.csv', index=False)
|
||||
pd.DataFrame(orders_data).to_csv('/opt/airflow/data/input/orders.csv', index=False)
|
||||
|
||||
print("Созданы файлы с данными клиентов и заказов")
|
||||
return "Sample data created"
|
||||
|
||||
def extract_customers():
|
||||
"""Извлечение данных клиентов"""
|
||||
df = pd.read_csv('/opt/airflow/data/input/customers.csv')
|
||||
print(f"Извлечено {len(df)} записей клиентов")
|
||||
|
||||
# Сохраняем извлеченные данные
|
||||
df.to_csv('/opt/airflow/data/output/extracted_customers.csv', index=False)
|
||||
return f"Извлечено {len(df)} клиентов"
|
||||
|
||||
def extract_orders():
|
||||
"""Извлечение данных заказов"""
|
||||
df = pd.read_csv('/opt/airflow/data/input/orders.csv')
|
||||
print(f"Извлечено {len(df)} записей заказов")
|
||||
|
||||
# Сохраняем извлеченные данные
|
||||
df.to_csv('/opt/airflow/data/output/extracted_orders.csv', index=False)
|
||||
return f"Извлечено {len(df)} заказов"
|
||||
|
||||
def transform_data():
|
||||
"""Преобразование данных - объединение клиентов и заказов"""
|
||||
customers_df = pd.read_csv('/opt/airflow/data/input/customers.csv')
|
||||
orders_df = pd.read_csv('/opt/airflow/data/input/orders.csv')
|
||||
|
||||
# Объединяем данные
|
||||
merged_df = pd.merge(orders_df, customers_df, on='customer_id', how='left')
|
||||
|
||||
# Добавляем вычисляемые поля
|
||||
merged_df['total_spent'] = merged_df['amount']
|
||||
merged_df['order_month'] = pd.to_datetime(merged_df['order_date']).dt.month
|
||||
|
||||
# Сохраняем преобразованные данные
|
||||
merged_df.to_csv('/opt/airflow/data/output/transformed_data.csv', index=False)
|
||||
print(f"Преобразованы данные: {len(merged_df)} записей")
|
||||
|
||||
return f"Преобразованы {len(merged_df)} записей"
|
||||
|
||||
def load_to_database():
|
||||
"""Загрузка данных в базу данных (симуляция)"""
|
||||
df = pd.read_csv('/opt/airflow/data/output/transformed_data.csv')
|
||||
|
||||
# В реальном сценарии здесь был бы код для загрузки в базу данных
|
||||
# Для учебных целей просто логируем
|
||||
print(f"Загружено в базу данных: {len(df)} записей")
|
||||
|
||||
# Создаем SQL для создания таблицы (в реальном сценарии)
|
||||
create_table_sql = """
|
||||
CREATE TABLE IF NOT EXISTS customer_orders (
|
||||
order_id INTEGER,
|
||||
customer_id INTEGER,
|
||||
product VARCHAR(100),
|
||||
amount DECIMAL(10,2),
|
||||
order_date DATE,
|
||||
name VARCHAR(100),
|
||||
email VARCHAR(100),
|
||||
join_date DATE,
|
||||
total_spent DECIMAL(10,2),
|
||||
order_month INTEGER
|
||||
);
|
||||
"""
|
||||
|
||||
print("SQL для создания таблицы:")
|
||||
print(create_table_sql)
|
||||
|
||||
return f"Подготовлено к загрузке в базу: {len(df)} записей"
|
||||
|
||||
def generate_report():
|
||||
"""Генерация отчета"""
|
||||
df = pd.read_csv('/opt/airflow/data/output/transformed_data.csv')
|
||||
|
||||
# Создаем простой отчет
|
||||
report = {
|
||||
'total_orders': len(df),
|
||||
'total_revenue': df['amount'].sum(),
|
||||
'avg_order_value': df['amount'].mean(),
|
||||
'unique_customers': df['customer_id'].nunique(),
|
||||
'top_customer': df.groupby('name')['amount'].sum().idxmax(),
|
||||
'top_customer_spending': df.groupby('name')['amount'].sum().max()
|
||||
}
|
||||
|
||||
# Сохраняем отчет в файл
|
||||
with open('/opt/airflow/data/output/report.txt', 'w') as f:
|
||||
f.write("Отчет по заказам клиентов\n")
|
||||
f.write("=" * 30 + "\n")
|
||||
f.write(f"Всего заказов: {report['total_orders']}\n")
|
||||
f.write(f"Общая выручка: ${report['total_revenue']}\n")
|
||||
f.write(f"Средний чек: ${report['avg_order_value']:.2f}\n")
|
||||
f.write(f"Уникальных клиентов: {report['unique_customers']}\n")
|
||||
f.write(f"Лучший клиент: {report['top_customer']} (${report['top_customer_spending']})\n")
|
||||
|
||||
print("Создан отчет по заказам")
|
||||
return "Отчет создан"
|
||||
|
||||
# Определение задач
|
||||
create_data_task = PythonOperator(
|
||||
task_id='create_sample_data',
|
||||
python_callable=create_sample_data,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
extract_customers_task = PythonOperator(
|
||||
task_id='extract_customers',
|
||||
python_callable=extract_customers,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
extract_orders_task = PythonOperator(
|
||||
task_id='extract_orders',
|
||||
python_callable=extract_orders,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
transform_task = PythonOperator(
|
||||
task_id='transform_data',
|
||||
python_callable=transform_data,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
load_task = PythonOperator(
|
||||
task_id='load_to_database',
|
||||
python_callable=load_to_database,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
report_task = PythonOperator(
|
||||
task_id='generate_report',
|
||||
python_callable=generate_report,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
# Установка зависимостей
|
||||
create_data_task >> [extract_customers_task, extract_orders_task] >> transform_task >> load_task >> report_task
|
||||
@@ -0,0 +1,107 @@
|
||||
"""
|
||||
DAG для демонстрации обработки ошибок в Airflow
|
||||
Уровень: Продвинутый
|
||||
"""
|
||||
from datetime import datetime, timedelta
|
||||
from airflow import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.operators.dummy import DummyOperator
|
||||
import random
|
||||
|
||||
# Определение DAG
|
||||
default_args = {
|
||||
'owner': 'student',
|
||||
'depends_on_past': False,
|
||||
'start_date': datetime(2023, 1, 1),
|
||||
'email_on_failure': False, # Отключаем email уведомления для простоты
|
||||
'email_on_retry': False,
|
||||
'retries': 3, # Количество попыток при ошибке
|
||||
'retry_delay': timedelta(seconds=10) # Задержка между попытками
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
'error_handling_dag',
|
||||
default_args=default_args,
|
||||
description='DAG для изучения обработки ошибок в Airflow',
|
||||
schedule_interval=timedelta(days=1),
|
||||
catchup=False,
|
||||
tags=['educational', 'error_handling', 'advanced']
|
||||
)
|
||||
|
||||
def unreliable_task():
|
||||
"""Задача, которая может завершиться с ошибкой"""
|
||||
# В реальном сценарии это может быть задача, зависящая от внешних факторов
|
||||
# Для учебных целей случайным образом генерируем ошибку
|
||||
if random.random() < 0.3: # 30% вероятность ошибки
|
||||
print("Ошибка: задача не выполнена успешно!")
|
||||
raise Exception("Случайная ошибка в задаче")
|
||||
|
||||
print("Задача выполнена успешно!")
|
||||
return "Задача выполнена"
|
||||
|
||||
def success_handler():
|
||||
"""Обработчик успешного выполнения"""
|
||||
print("Поздравляем! Все задачи выполнены успешно!")
|
||||
return "Успешно завершено"
|
||||
|
||||
def failure_handler():
|
||||
"""Обработчик ошибок"""
|
||||
print("Одна или несколько задач завершились с ошибкой!")
|
||||
print("Проверьте логи для получения дополнительной информации")
|
||||
return "Ошибка обработана"
|
||||
|
||||
def retry_task():
|
||||
"""Задача с механизмом повторных попыток"""
|
||||
# Имитируем задачу, которая может завершиться с ошибкой, но со временем исправляется
|
||||
import time
|
||||
time.sleep(2) # Имитация работы
|
||||
|
||||
# С вероятностью 50% задача завершится с ошибкой
|
||||
if random.random() < 0.5:
|
||||
print("Ошибка в retry_task!")
|
||||
raise Exception("Ошибка в задаче с повторными попытками")
|
||||
|
||||
print("retry_task выполнена успешно!")
|
||||
return "retry_task завершена"
|
||||
|
||||
# Определение задач
|
||||
start_task = DummyOperator(
|
||||
task_id='start_task',
|
||||
dag=dag
|
||||
)
|
||||
|
||||
unreliable_task = PythonOperator(
|
||||
task_id='unreliable_task',
|
||||
python_callable=unreliable_task,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
retry_task = PythonOperator(
|
||||
task_id='retry_task',
|
||||
python_callable=retry_task,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
success_handler_task = PythonOperator(
|
||||
task_id='success_handler',
|
||||
python_callable=success_handler,
|
||||
trigger_rule='all_success', # Выполняется только если все предыдущие задачи успешны
|
||||
dag=dag
|
||||
)
|
||||
|
||||
failure_handler_task = PythonOperator(
|
||||
task_id='failure_handler',
|
||||
python_callable=failure_handler,
|
||||
trigger_rule='one_failed', # Выполняется если хотя бы одна предыдущая задача завершилась с ошибкой
|
||||
dag=dag
|
||||
)
|
||||
|
||||
end_task = DummyOperator(
|
||||
task_id='end_task',
|
||||
dag=dag
|
||||
)
|
||||
|
||||
# Установка зависимостей
|
||||
start_task >> [unreliable_task, retry_task]
|
||||
[unreliable_task, retry_task] >> [success_handler_task, failure_handler_task]
|
||||
[success_handler_task, failure_handler_task] >> end_task
|
||||
@@ -0,0 +1,125 @@
|
||||
"""
|
||||
DAG для демонстрации работы с файлами в Airflow
|
||||
Уровень: Средний
|
||||
"""
|
||||
import pandas as pd
|
||||
import os
|
||||
from datetime import datetime, timedelta
|
||||
from airflow import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.operators.bash import BashOperator
|
||||
|
||||
# Определение DAG
|
||||
default_args = {
|
||||
'owner': 'student',
|
||||
'depends_on_past': False,
|
||||
'start_date': datetime(2023, 1, 1),
|
||||
'email_on_failure': False,
|
||||
'email_on_retry': False,
|
||||
'retries': 1,
|
||||
'retry_delay': timedelta(minutes=5)
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
'file_operations_dag',
|
||||
default_args=default_args,
|
||||
description='DAG для изучения работы с файлами в Airflow',
|
||||
schedule_interval=timedelta(days=1),
|
||||
catchup=False,
|
||||
tags=['educational', 'files', 'intermediate']
|
||||
)
|
||||
|
||||
def generate_sample_data():
|
||||
"""Создание CSV файла с примерными данными"""
|
||||
import pandas as pd
|
||||
import random
|
||||
|
||||
data = {
|
||||
'id': range(1, 101),
|
||||
'name': [f'User_{i}' for i in range(1, 101)],
|
||||
'age': [random.randint(18, 65) for _ in range(100)],
|
||||
'salary': [random.randint(30000, 100000) for _ in range(100)]
|
||||
}
|
||||
|
||||
df = pd.DataFrame(data)
|
||||
df.to_csv('/opt/airflow/data/input/sample_data.csv', index=False)
|
||||
print(f"Создан файл с {len(df)} записями")
|
||||
return "sample_data.csv created"
|
||||
|
||||
def read_and_validate_data():
|
||||
"""Чтение и валидация CSV файла"""
|
||||
df = pd.read_csv('/opt/airflow/data/input/sample_data.csv')
|
||||
print(f"Прочитан файл: {len(df)} строк, {len(df.columns)} столбцов")
|
||||
|
||||
# Простая валидация
|
||||
assert len(df) > 0, "Файл пустой"
|
||||
assert 'name' in df.columns, "Отсутствует столбец name"
|
||||
|
||||
return f"Файл валидирован: {len(df)} записей"
|
||||
|
||||
def transform_data():
|
||||
"""Преобразование данных"""
|
||||
df = pd.read_csv('/opt/airflow/data/input/sample_data.csv')
|
||||
|
||||
# Простое преобразование - добавим столбец с категорией зарплаты
|
||||
df['salary_category'] = df['salary'].apply(
|
||||
lambda x: 'High' if x >= 70000 else 'Medium' if x >= 50000 else 'Low'
|
||||
)
|
||||
|
||||
# Сохраняем обработанные данные
|
||||
df.to_csv('/opt/airflow/data/output/processed_data.csv', index=False)
|
||||
print(f"Обработаны данные: {len(df)} записей")
|
||||
|
||||
return f"Данные обработаны: {len(df)} записей"
|
||||
|
||||
def write_summary():
|
||||
"""Создание сводки по обработанным данным"""
|
||||
df = pd.read_csv('/opt/airflow/data/output/processed_data.csv')
|
||||
|
||||
summary = {
|
||||
'total_records': len(df),
|
||||
'avg_salary': df['salary'].mean(),
|
||||
'high_salary_count': len(df[df['salary_category'] == 'High']),
|
||||
'medium_salary_count': len(df[df['salary_category'] == 'Medium']),
|
||||
'low_salary_count': len(df[df['salary_category'] == 'Low'])
|
||||
}
|
||||
|
||||
# Сохраняем сводку в текстовый файл
|
||||
with open('/opt/airflow/data/output/summary.txt', 'w') as f:
|
||||
f.write("Сводка по обработанным данным:\n")
|
||||
f.write(f"Всего записей: {summary['total_records']}\n")
|
||||
f.write(f"Средняя зарплата: {summary['avg_salary']:.2f}\n")
|
||||
f.write(f"Высокая зарплата: {summary['high_salary_count']}\n")
|
||||
f.write(f"Средняя зарплата: {summary['medium_salary_count']}\n")
|
||||
f.write(f"Низкая зарплата: {summary['low_salary_count']}\n")
|
||||
|
||||
print("Создана сводка по данным")
|
||||
return "Сводка создана"
|
||||
|
||||
# Определение задач
|
||||
generate_task = PythonOperator(
|
||||
task_id='generate_sample_data',
|
||||
python_callable=generate_sample_data,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
read_task = PythonOperator(
|
||||
task_id='read_and_validate_data',
|
||||
python_callable=read_and_validate_data,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
transform_task = PythonOperator(
|
||||
task_id='transform_data',
|
||||
python_callable=transform_data,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
summary_task = PythonOperator(
|
||||
task_id='write_summary',
|
||||
python_callable=write_summary,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
# Установка зависимостей
|
||||
generate_task >> read_task >> transform_task >> summary_task
|
||||
@@ -0,0 +1,63 @@
|
||||
"""
|
||||
Простой DAG для демонстрации основных концепций Airflow
|
||||
Уровень: Начальный
|
||||
"""
|
||||
from datetime import datetime, timedelta
|
||||
from airflow import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.operators.bash import BashOperator
|
||||
|
||||
# Определение DAG
|
||||
default_args = {
|
||||
'owner': 'student',
|
||||
'depends_on_past': False,
|
||||
'start_date': datetime(2023, 1, 1),
|
||||
'email_on_failure': False,
|
||||
'email_on_retry': False,
|
||||
'retries': 1,
|
||||
'retry_delay': timedelta(minutes=5)
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
'hello_world_dag',
|
||||
default_args=default_args,
|
||||
description='Простой DAG для изучения основ Airflow',
|
||||
schedule_interval=timedelta(days=1),
|
||||
catchup=False,
|
||||
tags=['educational', 'beginner']
|
||||
)
|
||||
|
||||
# Функции для задач
|
||||
def print_hello():
|
||||
print("Hello World from Airflow!")
|
||||
return 'Hello World!'
|
||||
|
||||
def print_date():
|
||||
print(f"Current date: {datetime.now()}")
|
||||
return f"Date: {datetime.now()}"
|
||||
|
||||
def print_goodbye():
|
||||
print("Goodbye from Airflow!")
|
||||
return 'Goodbye!'
|
||||
|
||||
# Определение задач
|
||||
start_task = PythonOperator(
|
||||
task_id='start_task',
|
||||
python_callable=print_hello,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
date_task = PythonOperator(
|
||||
task_id='date_task',
|
||||
python_callable=print_date,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
end_task = PythonOperator(
|
||||
task_id='end_task',
|
||||
python_callable=print_goodbye,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
# Установка зависимостей
|
||||
start_task >> date_task >> end_task
|
||||
@@ -0,0 +1,86 @@
|
||||
"""
|
||||
DAG для демонстрации работы с SQL в Airflow
|
||||
Уровень: Начальный-Средний
|
||||
"""
|
||||
from datetime import datetime, timedelta
|
||||
from airflow import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.providers.postgres.operators.postgres import PostgresOperator
|
||||
|
||||
# Определение DAG
|
||||
default_args = {
|
||||
'owner': 'student',
|
||||
'depends_on_past': False,
|
||||
'start_date': datetime(2023, 1, 1),
|
||||
'email_on_failure': False,
|
||||
'email_on_retry': False,
|
||||
'retries': 1,
|
||||
'retry_delay': timedelta(minutes=5)
|
||||
}
|
||||
|
||||
dag = DAG(
|
||||
'sql_basic_dag',
|
||||
default_args=default_args,
|
||||
description='DAG для изучения SQL операций в Airflow',
|
||||
schedule_interval=timedelta(days=1),
|
||||
catchup=False,
|
||||
tags=['educational', 'sql', 'beginner']
|
||||
)
|
||||
|
||||
# SQL команды
|
||||
create_table_sql = """
|
||||
CREATE TABLE IF NOT EXISTS students (
|
||||
id SERIAL PRIMARY KEY,
|
||||
name VARCHAR(100),
|
||||
age INTEGER,
|
||||
created_at TIMESTAMP DEFAULT NOW()
|
||||
);
|
||||
"""
|
||||
|
||||
insert_data_sql = """
|
||||
INSERT INTO students (name, age) VALUES
|
||||
('Alice', 22),
|
||||
('Bob', 24),
|
||||
('Charlie', 21)
|
||||
ON CONFLICT DO NOTHING;
|
||||
"""
|
||||
|
||||
query_data_sql = """
|
||||
SELECT * FROM students;
|
||||
"""
|
||||
|
||||
drop_table_sql = """
|
||||
-- DROP TABLE IF EXISTS students; -- Закомментировано для сохранения данных
|
||||
"""
|
||||
|
||||
# Определение задач
|
||||
create_table_task = PostgresOperator(
|
||||
task_id='create_table',
|
||||
postgres_conn_id='postgres_training', # Это соединение нужно будет создать вручную в Airflow UI
|
||||
sql=create_table_sql,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
insert_data_task = PostgresOperator(
|
||||
task_id='insert_data',
|
||||
postgres_conn_id='postgres_training',
|
||||
sql=insert_data_sql,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
query_data_task = PostgresOperator(
|
||||
task_id='query_data',
|
||||
postgres_conn_id='postgres_training',
|
||||
sql=query_data_sql,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
drop_table_task = PostgresOperator(
|
||||
task_id='drop_table',
|
||||
postgres_conn_id='postgres_training',
|
||||
sql=drop_table_sql,
|
||||
dag=dag
|
||||
)
|
||||
|
||||
# Установка зависимостей
|
||||
create_table_task >> insert_data_task >> query_data_task >> drop_table_task
|
||||
Reference in New Issue
Block a user