Skip to content

Commit dae445e

Browse files
committed
feat: resolve airflow concurrency, global error handling, and modularize ui css
1 parent c95fda6 commit dae445e

5 files changed

Lines changed: 121 additions & 104 deletions

File tree

‎README.md‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,8 @@ DataPrep/
7575
│ ├── raw/ ← Datos crudos de entrada
7676
│ ├── interim/ ← Datos temporales de ejecución en Airflow
7777
│ └── processed/ ← Datasets limpios generados
78+
├── assets/
79+
│ └── style.css ← Hoja de estilos de la interfaz Streamlit
7880
├── dags/
7981
│ └── data_pipeline_dag.py ← DAG de Apache Airflow
8082
├── src/
@@ -271,6 +273,14 @@ Después del pipeline:
271273

272274
## 📋 Changelog
273275

276+
### v1.2.0 — 2026-06-15
277+
278+
**Seguridad, Resiliencia y Modularidad:**
279+
280+
- **Concurrencia Segura en Airflow:** Resuelta la vulnerabilidad crítica que causaba corrupción de datos en ejecuciones concurrentes. Ahora cada tarea del DAG usa `context["run_id"]` para nombrar los archivos temporales de forma única.
281+
- **Manejo Global de Errores:** Envuelta la ejecución central de `main.py` en un bloque `try-except`, para que fallos inesperados se registren formalmente en los logs de error del sistema (con trazabilidad) y eviten la detención abrupta y fallos feos para el usuario de CLI.
282+
- **UI Limpia y Modularizada:** El dashboard de Streamlit ha sido refactorizado, moviendo las dependencias visuales de `app.py` al archivo `assets/style.css`.
283+
274284
### v1.1.0 — 2026-06-02
275285

276286
**Refactorización y Mejoras de Rendimiento:**

‎app.py‎

Lines changed: 3 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -33,52 +33,9 @@
3333
)
3434

3535
# ── Custom CSS ────────────────────────────────────────────────────────────────
36-
st.markdown(
37-
"""
38-
<style>
39-
@import url('https://fonts.googleapis.com/css2?family=Inter:wght@300;400;500;600;700&display=swap');
40-
html, body, [class*="css"] { font-family: 'Inter', sans-serif; }
41-
42-
/* KPI cards */
43-
.kpi-card {
44-
background: linear-gradient(135deg, #1e293b 0%, #0f172a 100%);
45-
border: 1px solid #334155;
46-
border-radius: 16px;
47-
padding: 1.2rem 1rem;
48-
text-align: center;
49-
margin: 0.25rem 0;
50-
}
51-
.kpi-val { font-size: 2.2rem; font-weight: 700; color: #818cf8; line-height: 1.1; }
52-
.kpi-label { font-size: 0.8rem; color: #94a3b8; margin-top: 0.25rem; }
53-
.kpi-delta { font-size: 0.8rem; margin-top: 0.2rem; }
54-
.delta-ok { color: #22c55e; }
55-
.delta-warn { color: #f59e0b; }
56-
.delta-bad { color: #ef4444; }
57-
58-
/* Section headers */
59-
.section-title {
60-
font-size: 1.1rem; font-weight: 600;
61-
border-left: 3px solid #6366f1;
62-
padding-left: 0.75rem;
63-
margin: 1.5rem 0 0.75rem;
64-
color: #e2e8f0;
65-
}
66-
67-
/* Alert boxes */
68-
.alert-warn {
69-
background: #1c1012; border-left: 3px solid #f59e0b;
70-
padding: 0.6rem 1rem; border-radius: 8px;
71-
font-size: 0.85rem; color: #fbbf24; margin: 0.3rem 0;
72-
}
73-
.alert-ok {
74-
background: #0a1f12; border-left: 3px solid #22c55e;
75-
padding: 0.6rem 1rem; border-radius: 8px;
76-
font-size: 0.85rem; color: #4ade80;
77-
}
78-
</style>
79-
""",
80-
unsafe_allow_html=True,
81-
)
36+
css_path = Path(__file__).parent / "assets" / "style.css"
37+
if css_path.exists():
38+
st.markdown(f"<style>{css_path.read_text(encoding='utf-8')}</style>", unsafe_allow_html=True)
8239

8340

8441
# ── Helpers ───────────────────────────────────────────────────────────────────

‎assets/style.css‎

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
1+
@import url('https://fonts.googleapis.com/css2?family=Inter:wght@300;400;500;600;700&display=swap');
2+
html, body, [class*="css"] { font-family: 'Inter', sans-serif; }
3+
4+
/* KPI cards */
5+
.kpi-card {
6+
background: linear-gradient(135deg, #1e293b 0%, #0f172a 100%);
7+
border: 1px solid #334155;
8+
border-radius: 16px;
9+
padding: 1.2rem 1rem;
10+
text-align: center;
11+
margin: 0.25rem 0;
12+
}
13+
.kpi-val { font-size: 2.2rem; font-weight: 700; color: #818cf8; line-height: 1.1; }
14+
.kpi-label { font-size: 0.8rem; color: #94a3b8; margin-top: 0.25rem; }
15+
.kpi-delta { font-size: 0.8rem; margin-top: 0.2rem; }
16+
.delta-ok { color: #22c55e; }
17+
.delta-warn { color: #f59e0b; }
18+
.delta-bad { color: #ef4444; }
19+
20+
/* Section headers */
21+
.section-title {
22+
font-size: 1.1rem; font-weight: 600;
23+
border-left: 3px solid #6366f1;
24+
padding-left: 0.75rem;
25+
margin: 1.5rem 0 0.75rem;
26+
color: #e2e8f0;
27+
}
28+
29+
/* Alert boxes */
30+
.alert-warn {
31+
background: #1c1012; border-left: 3px solid #f59e0b;
32+
padding: 0.6rem 1rem; border-radius: 8px;
33+
font-size: 0.85rem; color: #fbbf24; margin: 0.3rem 0;
34+
}
35+
.alert-ok {
36+
background: #0a1f12; border-left: 3px solid #22c55e;
37+
padding: 0.6rem 1rem; border-radius: 8px;
38+
font-size: 0.85rem; color: #4ade80;
39+
}

‎dags/data_pipeline_dag.py‎

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -64,7 +64,8 @@ def task_ingest(**context):
6464
df = load_csv(DEFAULT_INPUT_FILE)
6565
logger.info(f"Ingested {len(df)} rows")
6666

67-
out_path = DATA_INTERIM_DIR / "df_raw.pkl"
67+
run_id = context["run_id"]
68+
out_path = DATA_INTERIM_DIR / f"df_raw_{run_id}.pkl"
6869
df.to_pickle(out_path)
6970
context["ti"].xcom_push(key="df_path", value=str(out_path))
7071
logger.info("Data saved to interim and path pushed to XCom")
@@ -132,7 +133,8 @@ def task_clean(**context):
132133
df_clean = clean_data(df, config=CLEANING_CONFIG)
133134
logger.info(f"Cleaned: {len(df)} → {len(df_clean)} rows")
134135

135-
out_path = DATA_INTERIM_DIR / "df_clean.pkl"
136+
run_id = context["run_id"]
137+
out_path = DATA_INTERIM_DIR / f"df_clean_{run_id}.pkl"
136138
df_clean.to_pickle(out_path)
137139
context["ti"].xcom_push(key="df_path", value=str(out_path))
138140

@@ -152,7 +154,8 @@ def task_transform(**context):
152154
df_transformed = transform_data(df)
153155
logger.info(f"Transformed: {len(df_transformed.columns)} output columns")
154156

155-
out_path = DATA_INTERIM_DIR / "df_transform.pkl"
157+
run_id = context["run_id"]
158+
out_path = DATA_INTERIM_DIR / f"df_transform_{run_id}.pkl"
156159
df_transformed.to_pickle(out_path)
157160
context["ti"].xcom_push(key="df_path", value=str(out_path))
158161

@@ -213,7 +216,8 @@ def task_load(**context):
213216
logger.info(f"Quality report saved to {DEFAULT_REPORT_FILE}")
214217

215218
# Cleanup interim files
216-
for p in ["df_raw.pkl", "df_clean.pkl", "df_transform.pkl"]:
219+
run_id = context["run_id"]
220+
for p in [f"df_raw_{run_id}.pkl", f"df_clean_{run_id}.pkl", f"df_transform_{run_id}.pkl"]:
217221
file_to_del = DATA_INTERIM_DIR / p
218222
if file_to_del.exists():
219223
file_to_del.unlink()

‎main.py‎

Lines changed: 61 additions & 54 deletions
Original file line numberDiff line numberDiff line change
@@ -72,60 +72,67 @@ def run_pipeline(
7272
logger.error("Run 'python generate_dataset.py' first to create test data.")
7373
return False
7474

75-
# ── STEP 2: VALIDATE (before) ───────────────────────────────────────────
76-
logger.info("[2/5] VALIDATION (before cleaning)")
77-
before_report = validate_data(
78-
df_raw,
79-
null_threshold_pct=VALIDATION_CONFIG["null_threshold_pct"],
80-
duplicate_threshold_pct=VALIDATION_CONFIG["duplicate_threshold_pct"],
81-
iqr_factor=VALIDATION_CONFIG["outlier_iqr_factor"],
82-
)
83-
logger.info(
84-
f" → {before_report.total_rows} rows | {before_report.duplicate_rows} dups "
85-
f"| {sum(before_report.null_counts.values())} nulls "
86-
f"| {len(before_report.alerts)} alerts"
87-
)
88-
if before_report.alerts:
89-
for alert in before_report.alerts:
90-
logger.warning(f" {alert}")
91-
92-
# ── STEP 3: CLEAN ───────────────────────────────────────────────────────
93-
logger.info("[3/5] CLEANING")
94-
df_clean = clean_data(df_raw, config=CLEANING_CONFIG)
95-
96-
# ── STEP 4: TRANSFORM ───────────────────────────────────────────────────
97-
logger.info("[4/5] TRANSFORMATION")
98-
df_transformed = transform_data(df_clean)
99-
100-
# ── STEP 5: VALIDATE (after) ────────────────────────────────────────────
101-
after_report = validate_data(df_transformed)
102-
logger.info(
103-
f" → {after_report.total_rows} rows | {after_report.duplicate_rows} dups "
104-
f"| {sum(after_report.null_counts.values())} nulls"
105-
)
106-
107-
# ── SAVE OUTPUT ─────────────────────────────────────────────────────────
108-
logger.info("[5/5] LOAD — Saving results")
109-
output_path.parent.mkdir(parents=True, exist_ok=True)
110-
df_transformed.to_csv(output_path, index=False, encoding="utf-8")
111-
logger.info(f" → Clean dataset saved: {output_path}")
112-
113-
# ── GENERATE REPORT ─────────────────────────────────────────────────────
114-
report_path.parent.mkdir(parents=True, exist_ok=True)
115-
generate_quality_report(
116-
before_report=before_report,
117-
after_report=after_report,
118-
report_path=report_path,
119-
version=PIPELINE_VERSION,
120-
)
121-
122-
elapsed = time.time() - start_time
123-
logger.info("=" * 60)
124-
logger.info(f"Pipeline completed successfully in {elapsed:.2f}s")
125-
logger.info(f" Clean dataset : {output_path}")
126-
logger.info(f" Quality report: {report_path}")
127-
logger.info("=" * 60)
128-
return True
75+
try:
76+
# ── STEP 2: VALIDATE (before) ───────────────────────────────────────────
77+
logger.info("[2/5] VALIDATION (before cleaning)")
78+
before_report = validate_data(
79+
df_raw,
80+
null_threshold_pct=VALIDATION_CONFIG["null_threshold_pct"],
81+
duplicate_threshold_pct=VALIDATION_CONFIG["duplicate_threshold_pct"],
82+
iqr_factor=VALIDATION_CONFIG["outlier_iqr_factor"],
83+
)
84+
logger.info(
85+
f" → {before_report.total_rows} rows | {before_report.duplicate_rows} dups "
86+
f"| {sum(before_report.null_counts.values())} nulls "
87+
f"| {len(before_report.alerts)} alerts"
88+
)
89+
if before_report.alerts:
90+
for alert in before_report.alerts:
91+
logger.warning(f" {alert}")
92+
93+
# ── STEP 3: CLEAN ───────────────────────────────────────────────────────
94+
logger.info("[3/5] CLEANING")
95+
df_clean = clean_data(df_raw, config=CLEANING_CONFIG)
96+
97+
# ── STEP 4: TRANSFORM ───────────────────────────────────────────────────
98+
logger.info("[4/5] TRANSFORMATION")
99+
df_transformed = transform_data(df_clean)
100+
101+
# ── STEP 5: VALIDATE (after) ────────────────────────────────────────────
102+
after_report = validate_data(df_transformed)
103+
logger.info(
104+
f" → {after_report.total_rows} rows | {after_report.duplicate_rows} dups "
105+
f"| {sum(after_report.null_counts.values())} nulls"
106+
)
107+
108+
# ── SAVE OUTPUT ─────────────────────────────────────────────────────────
109+
logger.info("[5/5] LOAD — Saving results")
110+
output_path.parent.mkdir(parents=True, exist_ok=True)
111+
df_transformed.to_csv(output_path, index=False, encoding="utf-8")
112+
logger.info(f" → Clean dataset saved: {output_path}")
113+
114+
# ── GENERATE REPORT ─────────────────────────────────────────────────────
115+
report_path.parent.mkdir(parents=True, exist_ok=True)
116+
generate_quality_report(
117+
before_report=before_report,
118+
after_report=after_report,
119+
report_path=report_path,
120+
version=PIPELINE_VERSION,
121+
)
122+
123+
elapsed = time.time() - start_time
124+
logger.info("=" * 60)
125+
logger.info(f"Pipeline completed successfully in {elapsed:.2f}s")
126+
logger.info(f" Clean dataset : {output_path}")
127+
logger.info(f" Quality report: {report_path}")
128+
logger.info("=" * 60)
129+
return True
130+
131+
except Exception as e:
132+
logger.error("=" * 60)
133+
logger.error(f"Pipeline failed during execution: {e}", exc_info=True)
134+
logger.error("=" * 60)
135+
return False
129136

130137

131138
def parse_args() -> argparse.Namespace:

0 commit comments

Comments
 (0)