forked from MLakshmi-191203/Dynamic_etl_framework
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathschema_engine.py
More file actions
149 lines (114 loc) · 3.6 KB
/
Copy pathschema_engine.py
File metadata and controls
149 lines (114 loc) · 3.6 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
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
# schema_engine.py
# =========================
# 🔹 CSV SCHEMA EXTRACTOR
# =========================
def get_csv_schema(df):
schema = []
for i, col in enumerate(df.columns):
schema.append({
"column_name": col.lower(),
"ordinal_position": i + 1,
"data_type": str(df[col].dtype)
})
return schema
# =========================
# 🔹 MYSQL SCHEMA EXTRACTOR
# =========================
def get_mysql_schema(conn, table):
cur = conn.cursor()
cur.execute(f"""
SELECT column_name, ordinal_position, data_type
FROM information_schema.columns
WHERE table_name = '{table}'
""")
schema = []
for row in cur.fetchall():
schema.append({
"column_name": row[0].lower(),
"ordinal_position": row[1],
"data_type": row[2]
})
return schema
# =========================
# 🔹 GET EXISTING METADATA
# =========================
def get_existing_metadata(conn, source_system, source_table):
cur = conn.cursor()
cur.execute("""
SELECT column_name, data_type
FROM dyn_etl.metadata
WHERE source_system=%s AND source_table=%s
""", (source_system, source_table))
return {row[0]: row[1] for row in cur.fetchall()}
# =========================
# 🔹 DETECT SCHEMA CHANGES
# =========================
def detect_schema_changes(source_schema, existing_meta):
actions = []
source_cols = {col["column_name"]: col for col in source_schema}
# 🔹 ADD + MODIFY
for col_name, col_data in source_cols.items():
if col_name not in existing_meta:
actions.append({
"action_type": "ADD",
"column": col_data
})
elif existing_meta[col_name] != col_data["data_type"]:
actions.append({
"action_type": "MODIFY",
"column": col_data
})
# 🔹 REMOVE
for col_name in existing_meta:
if col_name not in source_cols:
actions.append({
"action_type": "REMOVE",
"column_name": col_name
})
return actions
# =========================
# 🔹 INSERT INTO METADATA
# =========================
def insert_metadata(conn, config, actions):
cur = conn.cursor()
for act in actions:
if act["action_type"] in ["ADD", "MODIFY"]:
col = act["column"]
cur.execute("""
INSERT INTO dyn_etl.metadata(
source_system,
source_table,
column_name,
ordinal_position,
data_type,
action_type,
old_column_name,
new_column_name
)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s)
""", (
config["source_system"],
config["source_table"],
col["column_name"],
col["ordinal_position"],
col["data_type"],
act["action_type"],
None,
None
))
elif act["action_type"] == "REMOVE":
cur.execute("""
INSERT INTO dyn_etl.metadata (
source_system,
source_table,
column_name,
action_type
)
VALUES (%s,%s,%s,%s)
""", (
config["source_system"],
config["source_table"],
act["column_name"],
"REMOVE"
))
conn.commit()