Files
skillfactor-pipeline/pipeline/p1_load.py

188 lines
7.2 KiB
Python

"""Phase 1 — load ESCO, O*NET and the ESCO/O*NET crosswalk into MSSQL.
Sources (cached in data/raw/, downloaded by p0_download.py):
- ESCO v1.2.1 classification, EN, CSV (esco_v1.2.1/)
- O*NET 30.3 text database (db_30_3_text/)
- ESCO-O*NET occupation crosswalk (esco_onet_crosswalk.csv, header at line 17)
Idempotent: every table is dropped and re-created on each run.
Evidence store is Microsoft SQL Server (project requirement — no SQLite).
"""
import os
import sys
import pandas as pd
sys.path.insert(0, os.path.dirname(__file__))
from db import connect
RAW = os.path.join(os.path.dirname(__file__), "..", "data", "raw")
ESCO = os.path.join(RAW, "esco_v1.2.1")
ONET = os.path.join(RAW, "db_30_3_text")
def load_df(cn, df: pd.DataFrame, table: str, columns: dict, indexes=()):
"""columns: {df_col: (sql_col, sql_type)}; drops+creates table, bulk inserts."""
df = df[list(columns)].copy()
# ESCO CSVs enthalten vereinzelt Dubletten -> auf PK-Spalte deduplizieren
pk_cols = [df_col for df_col, (_, t) in columns.items() if "PRIMARY KEY" in t]
if pk_cols:
df = df.drop_duplicates(subset=pk_cols, keep="first")
df = df.where(pd.notna(df), None)
cur = cn.cursor()
cur.execute(f"IF OBJECT_ID('{table}','U') IS NOT NULL DROP TABLE [{table}]")
cols_sql = ", ".join(f"[{sql_col}] {sql_type}" for sql_col, sql_type in columns.values())
cur.execute(f"CREATE TABLE [{table}] ({cols_sql})")
placeholders = ", ".join("?" for _ in columns)
col_names = ", ".join(f"[{c[0]}]" for c in columns.values())
rows = [tuple(None if v is None else str(v) for v in row) for row in df.itertuples(index=False)]
# fast_executemany + NVARCHAR(MAX) buffers poorly -> only use it for short-column tables
cur.fast_executemany = not any("MAX" in t for _, t in columns.values())
for i in range(0, len(rows), 5000):
cur.executemany(
f"INSERT INTO [{table}] ({col_names}) VALUES ({placeholders})",
rows[i : i + 5000],
)
for idx_col in indexes:
cur.execute(f"CREATE INDEX IX_{table}_{idx_col} ON [{table}] ([{idx_col}])")
cn.commit()
print(f"{table}: {len(rows)} rows")
return len(rows)
def read_onet(name: str) -> pd.DataFrame:
return pd.read_csv(os.path.join(ONET, name), sep="\t", dtype=str, encoding="utf-8")
def main():
cn = connect()
counts = {}
# --- ESCO ---
occ = pd.read_csv(os.path.join(ESCO, "occupations_en.csv"), dtype=str)
counts["esco_occupation"] = load_df(
cn, occ, "esco_occupation",
{
"conceptUri": ("concept_uri", "NVARCHAR(200) NOT NULL PRIMARY KEY"),
"iscoGroup": ("isco_group", "NVARCHAR(10)"),
"preferredLabel": ("preferred_label", "NVARCHAR(400)"),
"altLabels": ("alt_labels", "NVARCHAR(MAX)"),
"description": ("description", "NVARCHAR(MAX)"),
"definition": ("definition", "NVARCHAR(MAX)"),
"code": ("code", "NVARCHAR(20)"),
},
)
skills = pd.read_csv(os.path.join(ESCO, "skills_en.csv"), dtype=str)
counts["esco_skill"] = load_df(
cn, skills, "esco_skill",
{
"conceptUri": ("concept_uri", "NVARCHAR(200) NOT NULL PRIMARY KEY"),
"skillType": ("skill_type", "NVARCHAR(50)"),
"reuseLevel": ("reuse_level", "NVARCHAR(50)"),
"preferredLabel": ("preferred_label", "NVARCHAR(400)"),
"altLabels": ("alt_labels", "NVARCHAR(MAX)"),
"description": ("description", "NVARCHAR(MAX)"),
},
)
rel = pd.read_csv(os.path.join(ESCO, "occupationSkillRelations_en.csv"), dtype=str)
counts["esco_occ_skill"] = load_df(
cn, rel, "esco_occ_skill",
{
"occupationUri": ("occupation_uri", "NVARCHAR(200) NOT NULL"),
"relationType": ("relation_type", "NVARCHAR(20)"),
"skillType": ("skill_type", "NVARCHAR(20)"),
"skillUri": ("skill_uri", "NVARCHAR(200) NOT NULL"),
},
indexes=("occupation_uri", "skill_uri"),
)
isco = pd.read_csv(os.path.join(ESCO, "ISCOGroups_en.csv"), dtype=str)
counts["esco_isco_group"] = load_df(
cn, isco, "esco_isco_group",
{
"code": ("code", "NVARCHAR(10) NOT NULL"),
"preferredLabel": ("preferred_label", "NVARCHAR(400)"),
"description": ("description", "NVARCHAR(MAX)"),
},
indexes=("code",),
)
# --- O*NET (30.3: "Software Skills" = former Technology Skills) ---
counts["onet_occupation"] = load_df(
cn, read_onet("Occupation Data.txt"), "onet_occupation",
{
"O*NET-SOC Code": ("soc_code", "NVARCHAR(15) NOT NULL PRIMARY KEY"),
"Title": ("title", "NVARCHAR(300)"),
"Description": ("description", "NVARCHAR(MAX)"),
},
)
counts["onet_task"] = load_df(
cn, read_onet("Task Statements.txt"), "onet_task",
{
"O*NET-SOC Code": ("soc_code", "NVARCHAR(15) NOT NULL"),
"Task ID": ("task_id", "NVARCHAR(15) NOT NULL"),
"Task": ("task", "NVARCHAR(MAX)"),
"Task Type": ("task_type", "NVARCHAR(30)"),
},
indexes=("soc_code",),
)
counts["onet_software"] = load_df(
cn, read_onet("Software Skills.txt"), "onet_software",
{
"O*NET-SOC Code": ("soc_code", "NVARCHAR(15) NOT NULL"),
"Workplace Example": ("example", "NVARCHAR(300)"),
"Element ID": ("element_id", "NVARCHAR(30)"),
"Element Name": ("element_name", "NVARCHAR(300)"),
"Hot Technology": ("hot_technology", "NVARCHAR(5)"),
"In Demand": ("in_demand", "NVARCHAR(5)"),
},
indexes=("soc_code",),
)
counts["onet_task_dwa"] = load_df(
cn, read_onet("Tasks to DWAs.txt"), "onet_task_dwa",
{
"O*NET-SOC Code": ("soc_code", "NVARCHAR(15) NOT NULL"),
"Task ID": ("task_id", "NVARCHAR(15) NOT NULL"),
"DWA Element ID": ("dwa_id", "NVARCHAR(30) NOT NULL"),
},
indexes=("soc_code", "dwa_id"),
)
dwa = read_onet("GWAs to IWAs to DWAs.txt")[["DWA Element ID", "DWA Element Name"]].drop_duplicates()
counts["onet_dwa"] = load_df(
cn, dwa, "onet_dwa",
{
"DWA Element ID": ("dwa_id", "NVARCHAR(30) NOT NULL PRIMARY KEY"),
"DWA Element Name": ("dwa_name", "NVARCHAR(400)"),
},
)
# --- ESCO <-> O*NET crosswalk (metadata preamble: header on line 17) ---
cw = pd.read_csv(os.path.join(RAW, "esco_onet_crosswalk.csv"), dtype=str, skiprows=16)
cw = cw[cw["ESCO or ISCO URI"].str.contains("/esco/occupation/", na=False)]
counts["crosswalk_esco_onet"] = load_df(
cn, cw, "crosswalk_esco_onet",
{
"O*NET Id": ("onet_id", "NVARCHAR(15) NOT NULL"),
"O*NET Title": ("onet_title", "NVARCHAR(300)"),
"ESCO or ISCO URI": ("esco_uri", "NVARCHAR(200) NOT NULL"),
"ESCO or ISCO Title": ("esco_title", "NVARCHAR(400)"),
"Type of Match": ("match_type", "NVARCHAR(30)"),
},
indexes=("esco_uri", "onet_id"),
)
cn.close()
print("\nSmoke test summary:")
for k, v in counts.items():
print(f" {k}: {v}")
if __name__ == "__main__":
main()