From a2bbceaebc1c38bd41dfd3867ef0e73b3d43c14f Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 6 Dec 2025 20:33:36 +0300 Subject: [PATCH] =?UTF-8?q?=D0=94=D0=BE=D0=B1=D0=B0=D0=B2=D0=BB=D0=B5?= =?UTF-8?q?=D0=BD=D0=B0=20=D0=BF=D1=80=D0=BE=D0=BF=D1=83=D1=89=D0=B5=D0=BD?= =?UTF-8?q?=D0=BD=D0=B0=D1=8F=20=D0=B8=D0=BB=D0=BB=D1=8E=D1=81=D1=82=D1=80?= =?UTF-8?q?=D0=B0=D1=86=D0=B8=D1=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- 09 - Продвинутые возможности Airflow.md | 45 ++++++++++++++++++++++--- 1 file changed, 41 insertions(+), 4 deletions(-) diff --git a/09 - Продвинутые возможности Airflow.md b/09 - Продвинутые возможности Airflow.md index f1ad8ea..8e1443b 100644 --- a/09 - Продвинутые возможности Airflow.md +++ b/09 - Продвинутые возможности Airflow.md @@ -177,6 +177,43 @@ with TaskGroup("data_loading") as loading_group: Обратите внимание: при использовании TaskGroup последовательность задач указывается внутри группы после объявления всех задач, а в конце DAG описывается последовательность выполнения самих групп. +Ниже приведена упрощённая схема зависимостей между тремя группами задач: + +```mermaid +flowchart LR + + %% group1 + subgraph G1["group1"] + g1_t1["task1"] + g1_t2["task2"] + g1_t3["task3"] + + g1_t1 --> g1_t2 + g1_t1 --> g1_t3 + end + + %% group2 + subgraph G2["group2"] + g2_t1["task1"] + g2_t2["task2"] + + g2_t1 --> g2_t2 + end + + %% group3 + subgraph G3["group3"] + g3_t1["task1"] + g3_t2["task2"] + + g3_t1 --> g3_t2 + end + + %% зависимости между группами + g1_t2 --> g3_t1 + g1_t3 --> g3_t1 + g2_t2 --> g3_t1 +``` + Визуально в интерфейсе Airflow группы задач отображаются как один узел с небольшим индикатором. Клик по нему разворачивает или сворачивает вложенные задачи, что значительно улучшает восприятие DAG с большим количеством задач и связей, особенно когда в них десятки и сотни задач. TaskGroup — это удобный способ логической группировки задач, который помогает упростить код и представить сложные пайплайны более компактно. @@ -312,8 +349,8 @@ def get_file_path(file_name): return os.path.join(os.path.expanduser('~/data'), file_name) def load_customer_data(): - url = 'https://example.com/customer_data.csv' - df = pd.read_csv(url) + file_path = get_file_path('customer_data.csv') + df = pd.read_csv(file_path) engine = create_engine(DATABASE_URL) df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging') @@ -396,8 +433,8 @@ def get_file_path(file_name): return os.path.join(os.path.expanduser('~/data'), file_name) def load_customer_data(): - url = 'https://example.com/customer_data.csv' - df = pd.read_csv(url) + file_path = get_file_path('customer_data.csv') + df = pd.read_csv(file_path) engine = create_engine(DATABASE_URL) df.to_sql('customers', engine, index=False, if_exists='replace', schema='staging')