-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtest_integration.py
More file actions
94 lines (72 loc) · 2.69 KB
/
Copy pathtest_integration.py
File metadata and controls
94 lines (72 loc) · 2.69 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
"""
tests/test_integration.py
--------------------------
Integration tests for Airflow DAG tasks and XCom pipeline data flow.
"""
import sys
from pathlib import Path
import pandas as pd
import pytest
# Skip Airflow integration tests gracefully if apache-airflow is not installed locally
pytest.importorskip("airflow")
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from dags.data_pipeline_dag import (
task_clean,
task_ingest,
task_load,
task_transform,
task_validate,
)
class MockTaskInstance:
"""Mock Airflow TaskInstance for testing XCom push/pull."""
def __init__(self):
self.xcom_store = {}
def xcom_push(self, key, value):
self.xcom_store[key] = value
def xcom_pull(self, task_ids, key):
return self.xcom_store.get(key)
class TestAirflowDAGIntegration:
"""Integration test suite for Airflow task execution sequence."""
def test_full_dag_task_sequence_integration(self, tmp_path, monkeypatch):
# Create input file
input_csv = tmp_path / "ventas_input.csv"
df = pd.DataFrame(
{
"id_venta": [1, 2, 2],
"fecha": ["2024-01-01", "2024-01-02", "2024-01-02"],
"precio": [100.0, 50.0, 50.0],
"cantidad": [2, 1, 1],
}
)
df.to_csv(input_csv, index=False)
# Patch default settings
import dags.data_pipeline_dag as dag_mod
monkeypatch.setattr(dag_mod, "DEFAULT_INPUT_FILE", input_csv)
monkeypatch.setattr(dag_mod, "DEFAULT_OUTPUT_FILE", tmp_path / "out.csv")
monkeypatch.setattr(dag_mod, "DEFAULT_REPORT_FILE", tmp_path / "report.html")
ti = MockTaskInstance()
context = {"run_id": "test_integration_run_123", "ti": ti}
# 1. Ingest
task_ingest(**context)
raw_pkl = ti.xcom_pull("ingest_data", "df_path")
assert Path(raw_pkl).exists()
# 2. Validate
task_validate(**context)
before_report_json = ti.xcom_pull("validate_data", "before_report")
assert before_report_json is not None
# 3. Clean
task_clean(**context)
clean_pkl = ti.xcom_pull("clean_data", "df_path")
assert Path(clean_pkl).exists()
# 4. Transform
task_transform(**context)
transform_pkl = ti.xcom_pull("transform_data", "df_path")
assert Path(transform_pkl).exists()
# 5. Load
task_load(**context)
assert (tmp_path / "out.csv").exists()
assert (tmp_path / "report.html").exists()
# Check interim files cleanup
assert not Path(raw_pkl).exists()
assert not Path(clean_pkl).exists()
assert not Path(transform_pkl).exists()