From 14d16c5caeaf6958b7c8b3f72ace9a3bb520791f Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 21:52:31 +0300 Subject: [PATCH] feat(airflow): implement dag orchestration for ddl and etl Add comprehensive DAG implementation for ClickHouse schema initialization and ETL pipeline orchestration. The ddl_init_dag manages database schema creation across stg/ods/dds/dm layers with verification capabilities. The etl_pipeline_dag implements full ODS to DDS to DM transformation flow with data quality checks, branching logic for full/incremental loads, and timeout handling for data availability. Additional changes: - Upgrade Airflow from 2.9.3 to 2.10.5 - Fix ClickHouse connection to use native protocol port 9000 - Mount SQL directory in docker-compose for DAG execution - Update project requirements and documentation comments - Remove unused pandas dependency --- Dockerfile.airflow | 2 +- airflow/requirements.txt | 9 +- dags/__pycache__/ddl_init_dag.cpython-312.pyc | Bin 0 -> 6608 bytes .../etl_pipeline_dag.cpython-312.pyc | Bin 0 -> 10146 bytes dags/ddl_init_dag.py | 188 ++++++++++++ dags/etl_pipeline_dag.py | 286 ++++++++++++++++-- docker-compose.yml | 9 +- 7 files changed, 459 insertions(+), 35 deletions(-) create mode 100644 dags/__pycache__/ddl_init_dag.cpython-312.pyc create mode 100644 dags/__pycache__/etl_pipeline_dag.cpython-312.pyc create mode 100644 dags/ddl_init_dag.py diff --git a/Dockerfile.airflow b/Dockerfile.airflow index 0643813..b857115 100644 --- a/Dockerfile.airflow +++ b/Dockerfile.airflow @@ -1,5 +1,5 @@ # Use the official Airflow image as base -FROM apache/airflow:2.9.3 +FROM apache/airflow:2.10.5 # Set environment variables ENV AIRFLOW_HOME=/opt/airflow diff --git a/airflow/requirements.txt b/airflow/requirements.txt index cdeea91..feae4fb 100644 --- a/airflow/requirements.txt +++ b/airflow/requirements.txt @@ -1,10 +1,7 @@ -# Airflow requirements для ClickHouse DWH проекта +# Airflow requirements для учебного ETL-проекта -# Core database connector для metadata +# Metadata DB для Airflow psycopg2-binary==2.9.9 -# ClickHouse provider для ETL +# ClickHouse operator/hook для DAG'ов airflow-clickhouse-plugin==1.6.0 - -# Для работы с данными -pandas==2.1.4 diff --git a/dags/__pycache__/ddl_init_dag.cpython-312.pyc b/dags/__pycache__/ddl_init_dag.cpython-312.pyc new file mode 100644 index 0000000000000000000000000000000000000000..9f0a80167f174f31c285def229b95b3bef23c7b4 GIT binary patch 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 tuple[str, ...]: + """Читает SQL-файл и делит его на отдельные команды по ';'.""" + file_path = SQL_ROOT / relative_path + if not file_path.is_file(): + raise AirflowException(f"SQL-файл не найден: {file_path}") + + sql_text = file_path.read_text(encoding="utf-8") + statements: list[str] = [] + for segment in sql_text.split(";"): + # Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы. + no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S) + lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")] + cleaned = "\n".join(lines).strip() + if cleaned: + statements.append(cleaned) + + if not statements: + raise AirflowException(f"SQL-файл пустой: {file_path}") + return tuple(statements) + + +# ----------------------------------------------------------------------------- +# SQL-проверки +# ----------------------------------------------------------------------------- +SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok" + +SQL_VERIFY_SCHEMA = """ +SELECT + (SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw, + (SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event, + (SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'click') AS dds_click, + (SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event, + (SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched +""" + + +# ----------------------------------------------------------------------------- +# Управляющие функции +# ----------------------------------------------------------------------------- +def choose_ddl_mode(**context) -> str: + """Выбирает ветку выполнения: full DDL или только verify.""" + dag_run = context.get("dag_run") + conf = dag_run.conf if dag_run else {} + verify_only = bool(conf.get("verify_only", context["params"]["verify_only"])) + return "skip_ddl" if verify_only else "ddl_00_databases" + + +def assert_schema_ready(**context) -> None: + """Проверяет результат финальной SQL-проверки схемы.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="verify_schema_sql") + + if not result or not result[0] or len(result[0]) != 5: + raise AirflowException(f"Некорректный результат проверки схемы: {result}") + + if any(value == 0 for value in result[0]): + raise AirflowException( + "Схема применена не полностью. Проверьте таблицы/VIEW stg, ods, dds, dm." + ) + + +with DAG( + dag_id="ddl_init", + description="Инициализация схемы ClickHouse (stg/ods/dds/dm)", + default_args=default_args, + schedule=None, + start_date=datetime(2024, 1, 1), + catchup=False, + max_active_runs=1, + is_paused_upon_creation=True, + tags=["ddl", "bootstrap", "clickhouse"], + params={ + "verify_only": Param(False, type="boolean"), + }, +) as dag: + check_clickhouse = ClickHouseOperator( + task_id="check_clickhouse", + sql=SQL_CHECK_CLICKHOUSE, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + choose_mode = BranchPythonOperator( + task_id="choose_mode", + python_callable=choose_ddl_mode, + ) + + ddl_00_databases = ClickHouseOperator( + task_id="ddl_00_databases", + sql=load_sql_statements("ddl/00_databases.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_10_stg = ClickHouseOperator( + task_id="ddl_10_stg", + sql=load_sql_statements("ddl/stg/10_stg.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_20_ods = ClickHouseOperator( + task_id="ddl_20_ods", + sql=load_sql_statements("ddl/ods/20_ods.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_30_dds = ClickHouseOperator( + task_id="ddl_30_dds", + sql=load_sql_statements("ddl/dds/30_dds.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_40_dm = ClickHouseOperator( + task_id="ddl_40_dm", + sql=load_sql_statements("ddl/dm/40_dm.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + skip_ddl = EmptyOperator(task_id="skip_ddl") + + ddl_complete = EmptyOperator( + task_id="ddl_complete", + trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, + ) + + verify_schema_sql = ClickHouseOperator( + task_id="verify_schema_sql", + sql=SQL_VERIFY_SCHEMA, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + verify_schema = PythonOperator( + task_id="verify_schema", + python_callable=assert_schema_ready, + ) + + check_clickhouse >> choose_mode + choose_mode >> skip_ddl >> ddl_complete + choose_mode >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds >> ddl_40_dm >> ddl_complete + ddl_complete >> verify_schema_sql >> verify_schema diff --git a/dags/etl_pipeline_dag.py b/dags/etl_pipeline_dag.py index feb0b1c..586f0dc 100644 --- a/dags/etl_pipeline_dag.py +++ b/dags/etl_pipeline_dag.py @@ -1,46 +1,282 @@ """ -ETL Pipeline DAG для ClickHouse DWH +DAG ETL-процесса ODS -> DDS -> DM для учебного проекта. -Шаблон DAG для оркестрации пайплайна данных. -Полная реализация будет добавлена позже. - -Пайплайн: - 1. DDL - создание структуры БД - 2. Load - загрузка данных в Kafka - 3. Transform - batch трансформация ODS → DDS → DM +Принципы реализации: +- SQL выполняется явными task на ClickHouseOperator; +- SQL-файлы вызываются по фиксированным путям; +- Python используется только для управляющей логики (branch/wait/assert). """ +from __future__ import annotations + +import re +import time from datetime import datetime, timedelta -from airflow import DAG -from airflow.operators.bash import BashOperator -from airflow.operators.empty import EmptyOperator +from pathlib import Path +from airflow import DAG +from airflow.exceptions import AirflowException +from airflow.models.param import Param +from airflow.operators.empty import EmptyOperator +from airflow.operators.python import BranchPythonOperator, PythonOperator +from airflow.utils.task_group import TaskGroup +from airflow.utils.trigger_rule import TriggerRule +from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook +from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator + + +# ----------------------------------------------------------------------------- # Базовые настройки DAG +# ----------------------------------------------------------------------------- default_args = { "owner": "airflow", "depends_on_past": False, "email_on_failure": False, "email_on_retry": False, "retries": 1, - "retry_delay": timedelta(minutes=5), + "retry_delay": timedelta(minutes=2), } + +# ----------------------------------------------------------------------------- +# SQL-файлы проекта +# ----------------------------------------------------------------------------- +SQL_ROOT = Path(__file__).resolve().parents[1] / "sql" + + +def load_sql_statements(relative_path: str) -> tuple[str, ...]: + """Читает SQL-файл и делит его на отдельные команды по ';'.""" + file_path = SQL_ROOT / relative_path + if not file_path.is_file(): + raise AirflowException(f"SQL-файл не найден: {file_path}") + + sql_text = file_path.read_text(encoding="utf-8") + statements: list[str] = [] + for segment in sql_text.split(";"): + # Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы. + no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S) + lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")] + cleaned = "\n".join(lines).strip() + if cleaned: + statements.append(cleaned) + + if not statements: + raise AirflowException(f"SQL-файл пустой: {file_path}") + return tuple(statements) + + +# ----------------------------------------------------------------------------- +# SQL для проверок и технических шагов +# ----------------------------------------------------------------------------- +SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok" + +SQL_CHECK_SCHEMA_READY = """ +SELECT + (SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw, + (SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event, + (SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event, + (SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched +""" + +SQL_CHECK_ODS_QUALITY = """ +SELECT + count() AS total_rows, + countIf(length(parse_errors) > 0) AS rows_with_errors, + round(if(count() = 0, 0, countIf(length(parse_errors) > 0) / count() * 100), 2) AS error_pct +FROM ods.browser_event +""" + +SQL_TRUNCATE_DDS_CLICK = "TRUNCATE TABLE dds.click" +SQL_TRUNCATE_DDS_EVENT = "TRUNCATE TABLE dds.event" + +SQL_CHECK_DDS_INTEGRITY = """ +SELECT + countIf(click_id IS NOT NULL AND click_id NOT IN (SELECT click_id FROM dds.click)) AS orphan_events +FROM dds.event +""" + +SQL_VALIDATE_DM_SUMMARY = "SELECT count() AS dq_rows FROM dm.dq_summary" + + +# ----------------------------------------------------------------------------- +# Управляющие функции +# ----------------------------------------------------------------------------- +def assert_schema_ready(**context) -> None: + """Падает, если DDL не применён полностью.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="precheck.check_schema_ready_sql") + + if not result or not result[0] or len(result[0]) != 4: + raise AirflowException(f"Некорректный результат check_schema_ready_sql: {result}") + + if any(value == 0 for value in result[0]): + raise AirflowException( + "Схема не готова: сначала запустите DAG ddl_init, затем повторите etl_pipeline." + ) + + +def wait_for_ods_data(**context) -> None: + """ + Ожидает появления строк в ods.browser_event до заданного таймаута. + Таймаут берётся из dag_run.conf.wait_ods_timeout_sec или из params. + """ + dag_run = context.get("dag_run") + conf = dag_run.conf if dag_run else {} + timeout_sec = int(conf.get("wait_ods_timeout_sec", context["params"]["wait_ods_timeout_sec"])) + poll_interval_sec = 10 + + hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default") + started = time.monotonic() + + while True: + rows = hook.execute("SELECT count() FROM ods.browser_event") + count_rows = int(rows[0][0]) if rows else 0 + if count_rows > 0: + return + + elapsed = int(time.monotonic() - started) + if elapsed >= timeout_sec: + raise AirflowException( + f"Таймаут ожидания ODS истёк ({timeout_sec} сек). " + "Таблица ods.browser_event всё ещё пуста." + ) + + time.sleep(poll_interval_sec) + + +def choose_full_refresh(**context) -> str: + """Ветвление: делать TRUNCATE DDS или пропустить.""" + dag_run = context.get("dag_run") + conf = dag_run.conf if dag_run else {} + full_refresh = bool(conf.get("full_refresh", context["params"]["full_refresh"])) + return "transform.truncate_dds_click" if full_refresh else "transform.skip_truncate" + + +def assert_dm_summary_not_empty(**context) -> None: + """Проверяет, что dm.dq_summary заполнена после загрузки.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="transform.validate_dm_summary_sql") + + if not result or not result[0] or len(result[0]) != 1: + raise AirflowException(f"Некорректный результат validate_dm_summary_sql: {result}") + + dq_rows = int(result[0][0]) + if dq_rows <= 0: + raise AirflowException("dm.dq_summary пуста после load_dm_summary.") + + with DAG( dag_id="etl_pipeline", + description="ETL ODS -> DDS -> DM для demo-проекта", default_args=default_args, - description="ETL pipeline для ClickHouse DWH", - schedule=None, # Запуск только вручную (пока) + schedule=None, start_date=datetime(2024, 1, 1), catchup=False, - tags=["etl", "clickhouse", "dwh"], + max_active_runs=1, + is_paused_upon_creation=True, + tags=["etl", "clickhouse", "demo"], + params={ + "full_refresh": Param(True, type="boolean"), + "wait_ods_timeout_sec": Param(600, type="integer", minimum=30), + }, ) as dag: - - # TODO: добавить задачи пайплайна - # - ddl: создание структуры БД - # - load: загрузка данных в Kafka - # - transform: batch трансформация - - start = EmptyOperator(task_id="start") - end = EmptyOperator(task_id="end") - - start >> end + with TaskGroup(group_id="precheck") as precheck: + check_clickhouse = ClickHouseOperator( + task_id="check_clickhouse", + sql=SQL_CHECK_CLICKHOUSE, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + check_schema_ready_sql = ClickHouseOperator( + task_id="check_schema_ready_sql", + sql=SQL_CHECK_SCHEMA_READY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + check_schema_ready = PythonOperator( + task_id="check_schema_ready", + python_callable=assert_schema_ready, + ) + + check_clickhouse >> check_schema_ready_sql >> check_schema_ready + + with TaskGroup(group_id="transform") as transform: + wait_for_ods_data_task = PythonOperator( + task_id="wait_for_ods_data", + python_callable=wait_for_ods_data, + ) + + check_ods_quality = ClickHouseOperator( + task_id="check_ods_quality", + sql=SQL_CHECK_ODS_QUALITY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + choose_refresh_mode = BranchPythonOperator( + task_id="choose_refresh_mode", + python_callable=choose_full_refresh, + ) + + truncate_dds_click = ClickHouseOperator( + task_id="truncate_dds_click", + sql=SQL_TRUNCATE_DDS_CLICK, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + truncate_dds_event = ClickHouseOperator( + task_id="truncate_dds_event", + sql=SQL_TRUNCATE_DDS_EVENT, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + skip_truncate = EmptyOperator(task_id="skip_truncate") + + truncate_complete = EmptyOperator( + task_id="truncate_complete", + trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, + ) + + load_dds = ClickHouseOperator( + task_id="load_dds", + sql=load_sql_statements("dds/30_ods_to_dds.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + check_dds_integrity = ClickHouseOperator( + task_id="check_dds_integrity", + sql=SQL_CHECK_DDS_INTEGRITY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + load_dm_summary = ClickHouseOperator( + task_id="load_dm_summary", + sql=load_sql_statements("dm/40_dds_to_dm.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + validate_dm_summary_sql = ClickHouseOperator( + task_id="validate_dm_summary_sql", + sql=SQL_VALIDATE_DM_SUMMARY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + validate_dm_summary = PythonOperator( + task_id="validate_dm_summary", + python_callable=assert_dm_summary_not_empty, + ) + + wait_for_ods_data_task >> check_ods_quality >> choose_refresh_mode + choose_refresh_mode >> truncate_dds_click >> truncate_dds_event >> truncate_complete + choose_refresh_mode >> skip_truncate >> truncate_complete + truncate_complete >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary_sql >> validate_dm_summary + + precheck >> transform diff --git a/docker-compose.yml b/docker-compose.yml index 5da65f9..976d03f 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -4,7 +4,7 @@ x-airflow-env: &airflow-default-env AIRFLOW__CORE__EXECUTOR: LocalExecutor AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW_SECRET_KEY:-replace-me-with-random-string} # ClickHouse connection для ETL - AIRFLOW_CONN_CLICKHOUSE_DEFAULT: clickhouse://default:123456@clickhouse:8123/default + AIRFLOW_CONN_CLICKHOUSE_DEFAULT: clickhouse://default:123456@clickhouse:9000/default services: @@ -97,7 +97,7 @@ services: build: context: . dockerfile: Dockerfile.airflow - image: airflow-optimized:2.9.3 + image: airflow-optimized:2.10.5 environment: <<: *airflow-default-env command: > @@ -108,6 +108,7 @@ services: - "8080:8080" volumes: - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: - cs_dwh @@ -132,6 +133,7 @@ services: " volumes: - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: - cs_dwh @@ -153,13 +155,14 @@ services: <<: *airflow-default-env volumes: - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: - cs_dwh command: > bash -ceuo pipefail " mkdir -p /opt/airflow/data && - chmod -R 777 /opt/airflow/data || true && + chmod -R a+rX /opt/airflow/data || true && chown -R airflow:0 /opt/airflow/data || true && umask 000 && su -s /bin/bash airflow -c 'airflow db migrate' &&