From 869c189fe859be21068755dcb6a54c24b427fbc9 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 8 Feb 2026 19:34:10 +0300 Subject: [PATCH] refactor(airflow): move DAGs to airflow/dags and update paths - Why: - keep Airflow artifacts under a single airflow/ directory - align repository layout with intended project structure - What: - move dags/ to airflow/dags/ and update compose mounts - make SQL root resolution work in container and local runs - update DAG path references in README, AGENTS, ARCHITECTURE, and plans - remove tracked Python cache artifacts from old DAG location - Check: - airflow dags list - airflow dags list-import-errors - e2e success: ddl_init, kafka_load(limit=50), etl_pipeline --- AGENTS.md | 10 +++++----- README.md | 2 +- {dags => airflow/dags}/.gitkeep | 0 {dags => airflow/dags}/__init__.py | 0 {dags => airflow/dags}/ddl_init_dag.py | 14 +++++++++++++- {dags => airflow/dags}/etl_pipeline_dag.py | 14 +++++++++++++- {dags => airflow/dags}/kafka_load_dag.py | 0 {dags => airflow/dags}/utils/__init__.py | 0 {dags => airflow/dags}/utils/kafka_helpers.py | 0 dags/__pycache__/ddl_init_dag.cpython-312.pyc | Bin 6608 -> 0 bytes .../etl_pipeline_dag.cpython-312.pyc | Bin 10146 -> 0 bytes docker-compose.yml | 6 +++--- docs/ARCHITECTURE.md | 6 +++--- plans/airflow_dags_plan.md | 6 +++--- 14 files changed, 41 insertions(+), 17 deletions(-) rename {dags => airflow/dags}/.gitkeep (100%) rename {dags => airflow/dags}/__init__.py (100%) rename {dags => airflow/dags}/ddl_init_dag.py (92%) rename {dags => airflow/dags}/etl_pipeline_dag.py (95%) rename {dags => airflow/dags}/kafka_load_dag.py (100%) rename {dags => airflow/dags}/utils/__init__.py (100%) rename {dags => airflow/dags}/utils/kafka_helpers.py (100%) delete mode 100644 dags/__pycache__/ddl_init_dag.cpython-312.pyc delete mode 100644 dags/__pycache__/etl_pipeline_dag.cpython-312.pyc diff --git a/AGENTS.md b/AGENTS.md index d03e985..45d58e6 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -20,11 +20,11 @@ ## Ключевые артефакты ### Исполняемые файлы (текущая структура) -- `dags/` — Airflow DAGs для оркестрации ETL: - - `dags/ddl_init_dag.py` — инициализация схемы ClickHouse - - `dags/etl_pipeline_dag.py` — ETL процесс STG → ODS → DDS → DM - - `dags/kafka_load_dag.py` — загрузка данных в Kafka из JSONL - - `dags/utils/kafka_helpers.py` — helper-функции для работы с Kafka +- `airflow/dags/` — Airflow DAGs для оркестрации ETL: + - `airflow/dags/ddl_init_dag.py` — инициализация схемы ClickHouse + - `airflow/dags/etl_pipeline_dag.py` — ETL процесс STG → ODS → DDS → DM + - `airflow/dags/kafka_load_dag.py` — загрузка данных в Kafka из JSONL + - `airflow/dags/utils/kafka_helpers.py` — helper-функции для работы с Kafka - `sql/` — SQL по слоям: - `sql/ddl/00_databases.sql` — создание БД stg/ods/dds/dm - `sql/ddl/stg/10_stg.sql` — STG слой (Kafka Engine + MV) diff --git a/README.md b/README.md index 83a1ebf..cb84c58 100644 --- a/README.md +++ b/README.md @@ -115,7 +115,6 @@ flowchart LR ``` . -├── dags/ # Airflow DAGs для оркестрации ├── sql/ │ ├── ddl/ # DDL по слоям │ │ ├── 00_databases.sql @@ -128,6 +127,7 @@ flowchart LR │ └── dm/ # Batch SQL: DDS -> DM ├── scripts/ # Служебные shell-скрипты (legacy fallback, не основной путь) ├── airflow/ # Конфигурация Airflow +│ ├── dags/ # Airflow DAGs для оркестрации │ └── requirements.txt ├── docs/ # Документация │ └── ARCHITECTURE.md # Подробное описание слоёв diff --git a/dags/.gitkeep b/airflow/dags/.gitkeep similarity index 100% rename from dags/.gitkeep rename to airflow/dags/.gitkeep diff --git a/dags/__init__.py b/airflow/dags/__init__.py similarity index 100% rename from dags/__init__.py rename to airflow/dags/__init__.py diff --git a/dags/ddl_init_dag.py b/airflow/dags/ddl_init_dag.py similarity index 92% rename from dags/ddl_init_dag.py rename to airflow/dags/ddl_init_dag.py index 9529177..ff34b0a 100644 --- a/dags/ddl_init_dag.py +++ b/airflow/dags/ddl_init_dag.py @@ -38,7 +38,19 @@ default_args = { # ----------------------------------------------------------------------------- # SQL-файлы проекта # ----------------------------------------------------------------------------- -SQL_ROOT = Path(__file__).resolve().parents[1] / "sql" +def resolve_sql_root() -> Path: + """Определяет корень SQL для контейнера и локального запуска.""" + candidates = ( + Path(__file__).resolve().parents[1] / "sql", # /opt/airflow/sql в контейнере + Path(__file__).resolve().parents[2] / "sql", # /sql при локальном запуске + ) + for candidate in candidates: + if candidate.is_dir(): + return candidate + return candidates[0] + + +SQL_ROOT = resolve_sql_root() def load_sql_statements(relative_path: str) -> tuple[str, ...]: diff --git a/dags/etl_pipeline_dag.py b/airflow/dags/etl_pipeline_dag.py similarity index 95% rename from dags/etl_pipeline_dag.py rename to airflow/dags/etl_pipeline_dag.py index eaff6c2..ce43d2b 100644 --- a/dags/etl_pipeline_dag.py +++ b/airflow/dags/etl_pipeline_dag.py @@ -41,7 +41,19 @@ default_args = { # ----------------------------------------------------------------------------- # SQL-файлы проекта # ----------------------------------------------------------------------------- -SQL_ROOT = Path(__file__).resolve().parents[1] / "sql" +def resolve_sql_root() -> Path: + """Определяет корень SQL для контейнера и локального запуска.""" + candidates = ( + Path(__file__).resolve().parents[1] / "sql", # /opt/airflow/sql в контейнере + Path(__file__).resolve().parents[2] / "sql", # /sql при локальном запуске + ) + for candidate in candidates: + if candidate.is_dir(): + return candidate + return candidates[0] + + +SQL_ROOT = resolve_sql_root() def load_sql_statements(relative_path: str) -> tuple[str, ...]: diff --git a/dags/kafka_load_dag.py b/airflow/dags/kafka_load_dag.py similarity index 100% rename from dags/kafka_load_dag.py rename to airflow/dags/kafka_load_dag.py diff --git a/dags/utils/__init__.py b/airflow/dags/utils/__init__.py similarity index 100% rename from dags/utils/__init__.py rename to airflow/dags/utils/__init__.py diff --git a/dags/utils/kafka_helpers.py b/airflow/dags/utils/kafka_helpers.py similarity index 100% rename from dags/utils/kafka_helpers.py rename to airflow/dags/utils/kafka_helpers.py diff --git a/dags/__pycache__/ddl_init_dag.cpython-312.pyc b/dags/__pycache__/ddl_init_dag.cpython-312.pyc deleted file mode 100644 index 9f0a80167f174f31c285def229b95b3bef23c7b4..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 6608 zcmb_gYj70TmA*aG^B!p)8oiJNjU*%uLL)F3Y$4;wdRQ_bY=rG%g15JsZb>uhdAPd= zJ!;q*V{eQX62L!#@~SACty)rOs}v_B@{4#~Tl=W|*d9byLzk&I6>nAiPee{-!>>K3 zXL?2g+dtOZQ+@m1d(J)Q+{gLuIrHm^3MYZ*$m_31pKc)J_t-EVZ=pi1?=%x~mGC4? zc#1cPRNNFcQCK&N=CB#Qv`EM4FdesqEqcF2w8m{=n_jny_P8VL(CaqQ8Fz(UdfhI% z;}ziwz3vbzLzSS4Z!@{q)SJ1i`W}oq*W>JB)OSG2FVvqy`x2wxRDY%30J1I+Itl_|;IDx6E{M1U zUyE5b{fS-vZg+DvAJU zZ{bh|#iM1*>N~ktU~X9u2aG^Nqu&F1?o}9bNqrAR1NjD+ybOdhf)pK}Ws?bURu=+m z@fKQIuoK6jsGq?Lrh|lw3>b~{i!g80TU3Q}en(wqU}k>{!r#kX2suB+LldMl7nevR z6;6pJ6Y^-#q}h2+5tL|L&>UFuf~au7MnB6blkhde9>K@k7nR1vCmg5sf15xqV|NB~Y-Jeo)=f*iDJmgICokTef3q=W=7!{V_iPF6Ip5a*&IHpk(c zmISS`)FTOsG^^S06BT65rMFp_IybAC<(I`s5h-!Whz(vqDr(&i%C8bd@5CZQ5dF?l z87fvpvJwd@271bsipA*Xf%kjl2PVU_&*4~<|G>kW7R?6FS%tXE7I+h}sucGk)}Vzp&Hqn5E5{9mT& ztUY6wsf%x9%rU&A9(Xb(ofNtFJ(^_9n|8?g;0`Hp=NFz#X-~#r&J?_7tTw|kgirDg{*d;Hk36Phh5v*S@wknAu>b4UWmvx+dI zXcjpoMitEjPG*+Vr!{I!bLw4UGy(o)k-^VW(4}yaB2Pz^NsW#rqY2H*rO;_Li;_+O z0ywk^34BnAo)N%@fy;r~3ROF}WkDH=QG~dVP-M*}3lmt=e2FA`T1-Z!*hn&7pwXyY zvlL9QMMQy12)vA5#2C;001r#i&GRA~O+*zITA|deRwX83W~i{NN;&`ocfl+F35qQF z#9KH2M9$lh9m%@_tFFeJtMRsL>qnJa)$M&B1V3m{w;f%n998Y3AA9OnJH@5tpe2E^AdNBZg=nU z$>qc9w&9gZ&^hvv%X@jtimOqj8~^f05AjV4^3lb}F)9pdj1RU&!&Rzd^T}fX4d|kpBI8_C0z*Ewu zcYg(a{Db-;7~2%|0FBj~|DWHd=ObJSIppu}Xh+O4jW1R5Zl`CTU&-rro;%mI~e> z2hEy!LQph1l1z+i^yy?$4APnnS_s3^<6yG%Bo;kT$e0+R7fSn~qIn{d$)qfxU^brQ z1?eg19YelnpvaPsE2}R}y)iW(T&Zl$K7FU5`TT)Py?@tx&6FMZ$X|0=xjJ)YW})KR z^wLbO`)eOG<+_ip)IPi7e=d9YPJQe79+mD@clJWL4${b`y~glI-|L%D!M)}w+yOms z`rMt`82n^RQcVW-1xH9er4(oxea8OrmMK7JX=AWu&6qX^TS_U^;%%=HL&88QR5N5r zuJUOrhbZT$Sx3;BJ_&}XKP(5@#n>D=Za_xxDhA^i2f$UYX8?FXZ%!$`EK1i*N>uX| zPlXH)73c6C!W$hk5jM)p61u>rxTfh5H01}N7lMXdR=40C&|S=o9yaJpH1{x@VqRQk z_(yn zuK6A3AHCf7&7PdI>o58VKj-mL=@f9NKPn_U7+7H<1PPVym*|_&rH4y~zuv?M2R7|90Dl)EW$semeX91++!0~l>Wt|=dAe9yflyF7h)duFkr+wh-K-m^ z`O9@k6Ls}~;)Mz#aM3v>M5#^80L=-U0cP2Hgrm8NIE)w!J|o}=PfE{}SvdS6hDFj7 z&^s4^73?^&58T7XE%hV6LWs#zy&X~%k)Qfv#^un@6`V}YkaYZyH? zI>-+7jT{*qV2_TBve+DZzQ2EPY)ncZ8=eLUflb1QLq{7uG`Vl@g%yYMa!7g&2pAH` z18{AunJJ38XR=ZBpZtVs`8{#`mUP@Db^lFj?~?Ajr0p)*^II}G|ic)!(k&yJj}It)Ep_SWRnWjG_*#Z2=*7 zYJ+)ibKcG58@ltpEuT3obyk%Q+^Zy2?RjrgzOMZUIL_9rres*}vjG z@X!E%Yu?+allQH-A1#sn*6cGN?yFFJEqQN)F8Pra_ui6ZpEWxIlD$m};oF|w+2P++ zF?sLSd^NLL-J7fK1;Yc@>=7`$zBBKIWoasntcN*P!_aCFye(hX39QxD?9nwRu~g5y zZqY4y+H?7_9NjRVxkYzEL#=Ab(Jc#eIeO1hI7dJEK~Ii8sy+wRu=-cATXZZ>d)2^> p9KGXO(=B@cpC?SX!hhD+`JP_hm(jOA+}mF<wt$WF%@KRFC{RTD zmPm2b5pa;cHR6oA0xr_GMch$Oz!NPAl#q9O#2YOQl#+f?q%2w~Q6_eZ8c4L~jYDh)QiO}8br>FqK4RB@88v?7j@<2UT5oq8l18cadz*??4(8$#Unz&VgX6|vW7I1vL`KuxOLn^ms;{OT7RUwY5j~mwGhm1D%B8c4${~JG9x$sBa}Wy?oO$co%{9k-0k`Fn1fl*~K?6o7)ZW60VkLSq>qC3AB34Sb!y?5{vdUfz7O88GJI06!hL&^FMvYYdV6*;eZzZN<)5c#3@%viwu|Bl)WQGx;i`&IHi==@at2 z-)58F1X5SzkAT<}7<`F=*UR#}`~keaE6)?~l}=kLGxX#hMt(c}lKgWZ^#M#c4ct!v zf8cUjej5NU$sfvB7%3=DFu-q~*%1kcCi>$^k>8i#g`gA{wyJbm)6WBk_dvpzRFdz) z+_!<*%Q=F`g8`%~Am(xT3N9AMfhysGt9TAjPs<-7#V01Ev3QIDyyL2H=~ME%>E{aK zoj{2J_&#ocwlQQU=7B8i6EHprtA3Axh2cH{Usy@=VIdd`jkO&OhNZTkDDr~j^V`0{ z%kQHVdoUJ@OF=0dkBNgmonqmF5-)|LykfIn-Yk@(TR zV`T;Ke>Hghxkt zVSh5h!>glk4E^!=1VB6mNS>xSsLHkTD>VFHp{)c=b2J=FO1$VZDaQEG7%wPJj!*D0 zPGsXTHW3sh#l=U1;RwEtz%wcEiX#sZcuAO4%=ifNqGBiSEbL5hQqhawjf8Se26vP) z!yg!M`mP6qe<*eM&bDO1}KZ88ob3>e<&o7J~#5slV|etyg$*M$VKnNI7>xQR9xnnAW>w z)*Q&RF=ZUj1rORNKDCwgQm0KJkS?Tqb`9`0rA*@`i}BRHhCN<-{}^Y{jo;NUvfgWbm-FQegtA?Ofuhjs`o_Usj~-grrl9T^MtuWm#R9|Z z{(bv~6?0f*N5TRxQ`<%{T^S%iw3uikr^*w6p%gF4e_2!e=FZo4p6{CX%s>6{_Fwh>YwwMk0ohWSt>197zB^stovDB9x2}=7HJ=#% z-S%&`Ul-l%XVd+x{2lH${roFl*)=k2geL1LJLfy=o9oMX>Ti0Q)1KyxXWgt>v3s&b zWjBjfrHfY0jsB*naiJL4-f>YyWkl|-D+jOamh1X64&XWPsoiy^Dr2va4K;uK+DW-b zX%Vf;FRW{N4W?iC4bVv{e29$so>$D$TQoRCQ%w~VtNL9nOHDN0EdP6QSfuD5IoxC zLr#FfE&{gv|9QnZab;fo5x7|@H^9krBj7)SAGe&}B0ez2{@L`SOP6y*r3i#$LO3+W zbLx_~D0|*lr0TiU^B`pst(@4U9BL_}Hj4)|?0mBG2-I5En!r)6KM# zGeUG0)jc<&PK9V7#x;jWnlW)mUFMs63J zmS4LYWH)V76~p@v4DRR|?qh~~w(o(<1}Gs616Ld&lcv6R-(J9i$s{(+F#|));J#sI z@W7ru6XSfQa-xj%+FIAh2DG&q0tiTi=w2whdx11T79-;F9;&P37+m9 z)xWlY7ln`-94aIh5d0HQh?&O~aGl@@(`V(6i0j030Gup%v-GL-%YH>Cg{SHg0uNdc zn(&jB72y#Lvgphv(cZrW-V;*=I5JVdaPuq7r&Sr=0P5hN(=#L&!Au3c_Iqk-f{{H0 z^B4qxpI@1HC$l89#I4{`P&uakzgdrizq3>2#>(yp0#U~kMYM6LibdonjS1!rN0XtWiTRk`}24S zm?nCm`91YtGt_lk>)hJeO=o(3*qOGq{!zpYIN@~D$w}zQKTf{@+W;|2ZsQ@GA<+Wl zL7fa73xp()+hL*x7?7!voCpI7jE^oem@0Bygbl~S5*#F^asZY;R1fTJB!wqQa2Gpf z|DeyHSdW2wXA{XtMA3ukQ3DZuFgEGaE9Oubl89RjUWjarQnJEs z2m^pI27jUm%?$OK)pf>ua@EZCY_aq7p_7N?hThpjnc}{gouAv?b7kl0zpMPm(1pf% z_glVyh+UztH*O)0lb1&#G9Q)eVi`wVw#2_Q+hFg`KYU^JkGI^gwcaqd0@{L(D(=_a z)$5$L1xjbNE$AtW|1PklmS`{P4AI{AL09ogjIX4#w5B_FS%f0UB=8O_`^D%mHLWY? z7UU5t(b9p1(4)R{I&8JVNuGnAGeB-p&}GOYAWz8gn>I@M++;kTTNDJ*GJt85m4b|- zxGMjkh@J!|NF;LkfzT)I{7_-6kW8}Y={{_rb z&VZ4-PKpyBpuK|I4z&ujQdm!TA4(6*eJ%`h2q+SuJ|Ph(IFNz>R3k9ycn|dvOAlB& zA^d)_7|KAQA<_EPI@nb20#V7|$1o5QAihKw%7h+@K@h$@% zJ{T2M;)*#pSMjbiPcf;;U*TO)Jj5}MVYLM^?2#ZSCCpblV8Rv;gCsF9yQE@4LcyuE z%RJ!9?p?Kr$#!5TFI)H3{I;vRjP)qkZrtm$u+z=NuA|ig!|wAlWcp8W)lgrWU+EOy zhPh|a|D42!mMhKsA$ie{@{(e}?Ns#Pm?YpuQVeS4B^rmCOgt72DQ5l{AA;J2YP*bL zgy$1Jo3IOzg+16Dz-jhH=b*R~@dz4Bm;}Zik47gj(Kr=%v*eGY&>4}pQbH0pY&CIS_ z)%CNRPj~&G>w<2k_ft>l8SdQJ*)h3p)BMpZ$I_eo%FeHzwPo$4bMz0ovqhdSJ*BU5=VE7LnX<+U<>|8ajHhGPoV9z-OkB6G{?t>S z^=!Yk^;+u!ML$9BrQzY;OW)P&OB@T7zS!}Njj|S>ww<)eC2MY2)-F_2cBk+-taD|n z`wEQwWmw&mUVu1kh~kVxl+Tnzk>yvg!hqF5)pE#}JGJrxIMVbfCLicx%~jK2QkKo% zR&cKAQ+{&Wl%f%ErV#y%fWI6Q27$v=jReI#5pD-0XkZ{hkZrS=n5gk`d0|wX2q#z# z#itjBVYYgG4Tt0Lh|jR>Rtp>;kZ1xgW6mXq;3k75lE!y3-x|GZjOblHr*>w_3NzcAso`Mt<}eXzxFC127=2_8N5l zG6H2SNLB;zyo=JOzzI)O!45ADmGcEL2$d~9Yw{Ub2C1B5T^db3DTcQ&=@;SBTxGb$ z8AC2Q13!XcGWsD1*r4DgK0NG9^xe+KWTSDFfFBYBDZVEzS=qU7}P$&v9yuwq^(=I3m zj6y#&;F#3>YS|^N&c%%dk}Q~WmdC2N5Qd4b!`}nr$DJ5I?)nq)YHIX6^rdd{|F9Cb6hvS!Db)hALIZPp#gIG&O%PhtG%C6`($S{y4Uu}B5q zV;=bVh2x{~+{YU9d~i4B!@q~{ZW!WdJx~Hyf@v%9+KK_O6*HsT0;Yok zX4}Hk*r53rFxgPdSa6Pn4+|y);a3j<(;dZ=`~KwTwGU&W8o7D30X&?D=qE9+fQgV& zt_{U=pBJH`Lret;h%4?NRj-Huj))Slz^deNl2RCGliNtHu7KH-Ql*i|4KF%C0YB*| z>v9OH&sKx)L?k&Hj`=Zg%+)Fajw#PZ@tgdP2TJ9e`>_6+RU-M{a^P@hs(fDXan-X3;;Ur+BrrL+KtfSrBv zK+m3m;e#r=7T>TC->9@!fcHJsH#n@67sw#;z~FG-uKh@e3_Jxyd!Zu8_U>he4(#3A zv;QEe1rV(iFa#l{UD1Qt5>oj3*VxQpa~zxJpi#>13#C~+E#Y~@cmbR5WAh?5v(PBj zgZl>i*quECd-{6Wy#s?Rz78GOv7>KjNH~dTFJW^En;&BHGB&5Nc?=s20@X5)r3xx`E9E0Hr05W>i8Ws zaGUbqrqCqtyq;UtIj%Iw@PZxKY631=}xiPrJE_bQ$-cI&RAaR zzG+{bwy)0E8*bVkO4}d0(37_N=S$P}hi5EzO&)sHLOWH_EK}YEgQ?0OQ)LS-Qzbl0 z7hF`;MtILwZa}EPPz})Xg(9kYfQHW9B0FunTT?|h-t9Ef2kE;m8~rp5Tus)jy)x_S z$+}l%i`}=J5bwn+XJFp9PD^NLk-+=a>;6&r3bYtM(~t=HDyV0PXy z>WfU@lsZhh1*j6xRSPvh?v~Xv_xKHK(*or+4$v1KxmDeo^)_Z*YqE74vZb}(SdBF{ z*-~-GLDg@}dRwxt#{57%4AjCvZA;eMoULj5#;C6aP}Kq?7?pC(`m|-e{BZwuOaHe% zqEcf66%m1 zOoLTz1_1~q+`KxuenZ-_LEilMb<5*f3w$~2NLxDOO}o;TUGnZj*DZ$tI5(ZPw8-r} z*DXC+OXb`%Y0FyK*Oj((U2~=_`{beVv}JsOqT=*+-PaUl-LCugX&u@2g@BGS7N4&ZnR*Wue*JSFuW-21jg#Z^!0-nYBy0^&Up|Ie~t*4g^2^}qBz)UP*w>~%x` o kafka_load -> etl_pipeline - публикации строк из `.jsonl` (`1 строка = 1 message value`). ### Общие helper-функции -- `dags/utils/clickhouse_helpers.py`: +- `airflow/dags/utils/clickhouse_helpers.py`: - `execute_sql(sql: str) -> None` - `execute_sql_file(path: str) -> None` - `fetch_one(sql: str) -> tuple` -- `dags/utils/kafka_helpers.py`: +- `airflow/dags/utils/kafka_helpers.py`: - `prepare_topics(reset: bool) -> None` - `load_jsonl(file_path: str, topic: str, limit: int) -> int` - `check_kafka_ready() -> None` ## Структура файлов ```text -dags/ +airflow/dags/ ├── __init__.py ├── ddl_init_dag.py # отдельный DAG для DDL (обязателен) ├── kafka_load_dag.py # отдельный DAG для ingest в Kafka (обязателен)