forked from MLakshmi-191203/Dynamic_etl_framework
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmain.py
More file actions
98 lines (73 loc) · 2.64 KB
/
Copy pathmain.py
File metadata and controls
98 lines (73 loc) · 2.64 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
from datetime import datetime
from db import get_pg_conn, get_mysql_conn
from config_reader import get_active_configs
from source_loader import load_source
from target_loader import create_table_if_not_exists, upsert_data
from audit import log_run
# 🔥 Schema
from schema_validation import validate_or_capture_schema
from schema_engine import get_csv_schema, get_mysql_schema
def run_pipeline():
configs = get_active_configs()
print("🔍 Config:", configs)
for config in configs:
start_time = datetime.now()
pg_conn = get_pg_conn()
try:
print(f"🚀 Running pipeline for {config['source_table']}")
# 🔹 Load source data
df = load_source(config)
# =========================
# 🔥 SCHEMA EXTRACTION
# =========================
if config["source_system"] == "CSV":
source_schema = get_csv_schema(df)
elif config["source_system"] == "MYSQL":
mysql_conn = get_mysql_conn(config["source_database"])
source_schema = get_mysql_schema(mysql_conn, config["source_table"])
mysql_conn.close()
else:
raise Exception("Unsupported source system")
# =========================
# 🔥 SCHEMA VALIDATION / CAPTURE
# =========================
validate_or_capture_schema(pg_conn, config, source_schema)
# =========================
# 🔹 TARGET DETAILS
# =========================
target_table = f"{config['target_schema']}.{config['target_table']}"
pk = config["primary_key"]
# 🔹 Create table
create_table_if_not_exists(pg_conn, target_table, df, pk)
# 🔹 Load data
ins, upd, rej = upsert_data(pg_conn, target_table, df, pk)
# 🔹 Audit log
log_run(
pg_conn,
config,
start_time,
datetime.now(),
ins,
upd,
rej,
"SUCCESS",
None
)
print(f"✅ Completed {target_table}")
except Exception as e:
print(f"❌ Failed: {e}")
# 🔴 VERY IMPORTANT FIX
pg_conn.rollback()
log_run(
pg_conn,
config,
start_time,
datetime.now(),
0, 0, 0,
"FAILED",
str(e)
)
finally:
pg_conn.close()
if __name__ == "__main__":
run_pipeline()