forked from MLakshmi-191203/Dynamic_etl_framework
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtarget_loader.py
More file actions
59 lines (39 loc) · 1.21 KB
/
Copy pathtarget_loader.py
File metadata and controls
59 lines (39 loc) · 1.21 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
def create_table_if_not_exists(conn, table, df, pk):
cur = conn.cursor()
columns = []
for col in df.columns:
columns.append(f"{col} TEXT")
columns.append(f"PRIMARY KEY ({pk})")
ddl = f"""
CREATE TABLE IF NOT EXISTS {table} (
{",".join(columns)}
)
"""
cur.execute(ddl)
conn.commit()
def upsert_data(conn, table, df, pk):
cur = conn.cursor()
insert_count = 0
update_count = 0
reject_count = 0
for _, row in df.iterrows():
try:
cols = list(df.columns)
values = tuple(row)
col_str = ",".join(cols)
val_str = ",".join(["%s"] * len(values))
update_str = ",".join([f"{c}=EXCLUDED.{c}" for c in cols])
query = f"""
INSERT INTO {table} ({col_str})
VALUES ({val_str})
ON CONFLICT ({pk})
DO UPDATE SET {update_str}
"""
cur.execute(query, values)
# 👉 NOTE: we don’t know insert/update here
update_count += 1
except Exception as e:
print("❌ Reject:", e)
reject_count += 1
conn.commit()
return insert_count, update_count, reject_count