File tree Expand file tree Collapse file tree
Expand file tree Collapse file tree Original file line number Diff line number Diff line change 3434else :
3535 from airflow .providers .common .compat .sdk import DAG
3636
37- if AIRFLOW_V_3_0_PLUS :
37+ from airflow .serialization .serialized_objects import SerializedDAG
38+
39+ try :
3840 from airflow .serialization .serialized_objects import DagSerialization
39- else :
40- from airflow . serialization . serialized_objects import SerializedDAG
41+ except ImportError :
42+ DagSerialization = SerializedDAG
4143
4244DATA_INTERVAL_START = pendulum .datetime (2022 , 1 , 1 , tz = "UTC" )
4345DATA_INTERVAL_END = DATA_INTERVAL_START + dt .timedelta (hours = 1 )
@@ -62,20 +64,12 @@ def sync_dag_to_db(
6264
6365 def _write_dag (dag : DAG ) -> SerializedDAG :
6466 if not SerializedDagModel .has_dag (dag .dag_id ):
65- data = (
66- DagSerialization .to_dict (dag )
67- if AIRFLOW_V_3_0_PLUS
68- else SerializedDAG .to_dict (dag )
69- )
67+ data = DagSerialization .to_dict (dag )
7068 SerializedDagModel .write_dag (
7169 LazyDeserializedDAG (data = data ), bundle_name , session = session
7270 )
7371 session .flush ()
74- return (
75- DagSerialization .from_dict (data )
76- if AIRFLOW_V_3_0_PLUS
77- else SerializedDAG .from_dict (data )
78- )
72+ return DagSerialization .from_dict (data )
7973
8074 SerializedDAG .bulk_write_to_db (bundle_name , None , [dag ], session = session )
8175 _ = _write_dag (dag )
You can’t perform that action at this time.
0 commit comments