Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
56c6b79a5a | ||
|
|
452cfdef85 | ||
|
|
96dd1f834c | ||
|
|
e66ef86729 | ||
|
|
0318ee6cc5 | ||
|
|
0ddfb1997d |
@@ -17,3 +17,9 @@ APP_LOG_LEVEL=INFO
|
||||
APP_SECRET_KEY=change_me
|
||||
APP_CORS_ORIGINS=http://localhost:4200
|
||||
BACKEND_PORT=8000
|
||||
|
||||
# API Mock EnerVision
|
||||
APP_MOCK_API_BASE_URL=https://api-mock.charlieandre.fr
|
||||
APP_MOCK_API_USERNAME=change_me
|
||||
APP_MOCK_API_PASSWORD=change_me
|
||||
APP_MOCK_API_TIMEOUT_SECONDS=10
|
||||
|
||||
@@ -18,3 +18,7 @@ APP_SMTP_HOST=localhost
|
||||
APP_SMTP_PORT=1025
|
||||
APP_SMTP_USE_TLS=false
|
||||
APP_SMTP_FROM_ADDRESS=no-reply@enervision.fr
|
||||
APP_MOCK_API_BASE_URL=https://api-mock.charlieandre.fr
|
||||
APP_MOCK_API_USERNAME=change_me
|
||||
APP_MOCK_API_PASSWORD=change_me
|
||||
APP_MOCK_API_TIMEOUT_SECONDS=10
|
||||
|
||||
@@ -34,6 +34,11 @@ class Settings(BaseSettings):
|
||||
database_pool_size: int = 5
|
||||
database_max_overflow: int = 10
|
||||
|
||||
mock_api_base_url: str = "https://api-mock.charlieandre.fr"
|
||||
mock_api_username: str | None = None
|
||||
mock_api_password: SecretStr | None = None
|
||||
mock_api_timeout_seconds: float = Field(default=10.0, gt=0)
|
||||
|
||||
jwt_issuer: str = "enervision-api"
|
||||
jwt_audience: str = "enervision-web"
|
||||
access_token_ttl_seconds: int = Field(default=900, ge=60, le=3600)
|
||||
|
||||
@@ -0,0 +1,315 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import asyncio
|
||||
import json
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.ext.asyncio import AsyncConnection, create_async_engine
|
||||
|
||||
from app.core.config import get_settings
|
||||
|
||||
SOURCE_HISTORY = "api_history"
|
||||
|
||||
|
||||
def create_mock_api_client() -> httpx.AsyncClient:
|
||||
settings = get_settings()
|
||||
|
||||
if settings.mock_api_username is None or settings.mock_api_password is None:
|
||||
raise ValueError("Les identifiants de l'API Mock ne sont pas configurés.")
|
||||
|
||||
return httpx.AsyncClient(
|
||||
base_url=settings.mock_api_base_url.rstrip("/"),
|
||||
auth=(
|
||||
settings.mock_api_username,
|
||||
settings.mock_api_password.get_secret_value(),
|
||||
),
|
||||
timeout=settings.mock_api_timeout_seconds,
|
||||
)
|
||||
|
||||
|
||||
async def fetch_sites(
|
||||
client: httpx.AsyncClient,
|
||||
) -> list[dict[str, Any]]:
|
||||
response = await client.get("/api/v1/sites")
|
||||
|
||||
response.raise_for_status()
|
||||
|
||||
payload = response.json()
|
||||
|
||||
if not isinstance(payload, list):
|
||||
raise ValueError("La réponse /api/v1/sites doit être une liste.")
|
||||
|
||||
return payload
|
||||
|
||||
|
||||
async def upsert_sites(
|
||||
connection: AsyncConnection,
|
||||
sites: list[dict[str, Any]],
|
||||
) -> None:
|
||||
if not sites:
|
||||
return
|
||||
|
||||
await connection.execute(
|
||||
text(
|
||||
"""
|
||||
INSERT INTO site (
|
||||
site_id,
|
||||
site_type,
|
||||
site_name,
|
||||
location,
|
||||
capacity_kw,
|
||||
status
|
||||
)
|
||||
VALUES (
|
||||
:site_id,
|
||||
:site_type,
|
||||
:site_name,
|
||||
:location,
|
||||
:capacity_kw,
|
||||
:status
|
||||
)
|
||||
ON CONFLICT (site_id)
|
||||
DO UPDATE SET
|
||||
site_type = EXCLUDED.site_type,
|
||||
site_name = EXCLUDED.site_name,
|
||||
location = EXCLUDED.location,
|
||||
capacity_kw = EXCLUDED.capacity_kw,
|
||||
status = EXCLUDED.status
|
||||
"""
|
||||
),
|
||||
sites,
|
||||
)
|
||||
|
||||
|
||||
async def fetch_readings(
|
||||
client: httpx.AsyncClient,
|
||||
site_id: str,
|
||||
start_time: datetime,
|
||||
end_time: datetime,
|
||||
limit: int = 1000,
|
||||
) -> list[dict[str, Any]]:
|
||||
response = await client.get(
|
||||
"/api/v1/readings",
|
||||
params={
|
||||
"site_id": site_id,
|
||||
"start_time": start_time.isoformat(),
|
||||
"end_time": end_time.isoformat(),
|
||||
"limit": limit,
|
||||
},
|
||||
)
|
||||
|
||||
response.raise_for_status()
|
||||
|
||||
payload = response.json()
|
||||
|
||||
if not isinstance(payload, list):
|
||||
raise ValueError("La réponse /api/v1/readings doit être une liste.")
|
||||
|
||||
return payload
|
||||
|
||||
|
||||
def build_reading_row(
|
||||
reading: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
timestamp = datetime.fromisoformat(reading["timestamp"].replace("Z", "+00:00"))
|
||||
return {
|
||||
"site_id": reading["site_id"],
|
||||
"timestamp": timestamp,
|
||||
"source": SOURCE_HISTORY,
|
||||
"dataset_id": None,
|
||||
"consumption_kw": reading.get("consumption_kw"),
|
||||
"consumption_kwh": reading.get("consumption_kwh"),
|
||||
"consumption_euros": None,
|
||||
"voltage_v": reading.get("voltage_v"),
|
||||
"current_a": reading.get("current_a"),
|
||||
"power_factor": reading.get("power_factor"),
|
||||
"temperature_celsius": reading.get("temperature_celsius"),
|
||||
"humidity_percent": reading.get("humidity_percent"),
|
||||
"solar_irradiance_wm2": None,
|
||||
"is_working_hours": None,
|
||||
"data_quality": reading.get("data_quality"),
|
||||
"null_reasons": reading.get("null_reasons"),
|
||||
"imputed_values": None,
|
||||
"imputation_method": None,
|
||||
"raw_data": json.dumps(
|
||||
reading,
|
||||
ensure_ascii=False,
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
READING_INSERT = text(
|
||||
"""
|
||||
INSERT INTO reading (
|
||||
site_id,
|
||||
timestamp,
|
||||
source,
|
||||
dataset_id,
|
||||
consumption_kw,
|
||||
consumption_kwh,
|
||||
consumption_euros,
|
||||
voltage_v,
|
||||
current_a,
|
||||
power_factor,
|
||||
temperature_celsius,
|
||||
humidity_percent,
|
||||
solar_irradiance_wm2,
|
||||
is_working_hours,
|
||||
data_quality,
|
||||
null_reasons,
|
||||
imputed_values,
|
||||
imputation_method,
|
||||
raw_data
|
||||
)
|
||||
VALUES (
|
||||
:site_id,
|
||||
:timestamp,
|
||||
:source,
|
||||
:dataset_id,
|
||||
:consumption_kw,
|
||||
:consumption_kwh,
|
||||
:consumption_euros,
|
||||
:voltage_v,
|
||||
:current_a,
|
||||
:power_factor,
|
||||
:temperature_celsius,
|
||||
:humidity_percent,
|
||||
:solar_irradiance_wm2,
|
||||
:is_working_hours,
|
||||
:data_quality,
|
||||
:null_reasons,
|
||||
CAST(:imputed_values AS jsonb),
|
||||
:imputation_method,
|
||||
CAST(:raw_data AS jsonb)
|
||||
)
|
||||
ON CONFLICT DO NOTHING
|
||||
"""
|
||||
)
|
||||
|
||||
|
||||
def build_reading_batch(
|
||||
readings: list[dict[str, Any]],
|
||||
) -> list[dict[str, Any]]:
|
||||
return [build_reading_row(reading) for reading in readings]
|
||||
|
||||
|
||||
async def import_mock_api_history(
|
||||
start_time: datetime,
|
||||
end_time: datetime,
|
||||
limit: int,
|
||||
dry_run: bool,
|
||||
) -> None:
|
||||
settings = get_settings()
|
||||
|
||||
async with create_mock_api_client() as client:
|
||||
sites = await fetch_sites(client)
|
||||
|
||||
print(f"Sites récupérés : {len(sites)}")
|
||||
|
||||
all_readings: list[dict[str, Any]] = []
|
||||
|
||||
for site in sites:
|
||||
site_id = site["site_id"]
|
||||
|
||||
readings = await fetch_readings(
|
||||
client=client,
|
||||
site_id=site_id,
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
limit=limit,
|
||||
)
|
||||
|
||||
print(f"{site_id}: {len(readings)} lectures")
|
||||
|
||||
all_readings.extend(readings)
|
||||
|
||||
print(f"Lectures récupérées : {len(all_readings)}")
|
||||
|
||||
if dry_run:
|
||||
print("Dry-run terminé : aucune donnée écrite.")
|
||||
return
|
||||
|
||||
engine = create_async_engine(
|
||||
str(settings.database_url),
|
||||
pool_pre_ping=True,
|
||||
)
|
||||
|
||||
try:
|
||||
async with engine.begin() as connection:
|
||||
await upsert_sites(
|
||||
connection,
|
||||
sites,
|
||||
)
|
||||
|
||||
rows = build_reading_batch(all_readings)
|
||||
|
||||
if rows:
|
||||
await connection.execute(
|
||||
READING_INSERT,
|
||||
rows,
|
||||
)
|
||||
|
||||
finally:
|
||||
await engine.dispose()
|
||||
|
||||
print("Import API Mock terminé.")
|
||||
|
||||
|
||||
def parse_datetime(value: str) -> datetime:
|
||||
return datetime.fromisoformat(value.replace("Z", "+00:00"))
|
||||
|
||||
|
||||
def parse_args() -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(description=("Import historique depuis l'API Mock EnerVision"))
|
||||
|
||||
parser.add_argument(
|
||||
"--start-time",
|
||||
required=True,
|
||||
type=parse_datetime,
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--end-time",
|
||||
required=True,
|
||||
type=parse_datetime,
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--limit",
|
||||
type=int,
|
||||
default=1000,
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--dry-run",
|
||||
action="store_true",
|
||||
)
|
||||
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
args = parse_args()
|
||||
|
||||
if args.limit < 1 or args.limit > 1000:
|
||||
raise ValueError("--limit doit être compris entre 1 et 1000.")
|
||||
|
||||
if args.start_time >= args.end_time:
|
||||
raise ValueError("--start-time doit être antérieur à --end-time.")
|
||||
|
||||
asyncio.run(
|
||||
import_mock_api_history(
|
||||
start_time=args.start_time,
|
||||
end_time=args.end_time,
|
||||
limit=args.limit,
|
||||
dry_run=args.dry_run,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -17,6 +17,7 @@ dependencies = [
|
||||
"argon2-cffi>=23.1",
|
||||
"anyio>=4.0",
|
||||
"aiosmtplib>=5.1.3",
|
||||
"httpx>=0.28.1",
|
||||
"pandas>=3.0.5",
|
||||
]
|
||||
|
||||
@@ -27,7 +28,6 @@ dev = [
|
||||
"pytest>=9.1.1",
|
||||
"pytest-asyncio>=1.4.0",
|
||||
"pytest-cov>=7.1.0",
|
||||
"httpx>=0.28.1",
|
||||
"pandas-stubs>=3.0.5.260914",
|
||||
]
|
||||
|
||||
|
||||
@@ -0,0 +1,622 @@
|
||||
import json
|
||||
import sys
|
||||
from datetime import datetime
|
||||
from types import SimpleNamespace
|
||||
from typing import Any
|
||||
from unittest.mock import AsyncMock, MagicMock
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
from httpx import AsyncClient, MockTransport, Request, Response
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
import app.etl.mock_api_import as mock_api_import
|
||||
from app.etl.mock_api_import import (
|
||||
READING_INSERT,
|
||||
SOURCE_HISTORY,
|
||||
build_reading_batch,
|
||||
build_reading_row,
|
||||
fetch_readings,
|
||||
fetch_sites,
|
||||
upsert_sites,
|
||||
)
|
||||
|
||||
|
||||
def make_site() -> dict[str, Any]:
|
||||
return {
|
||||
"site_id": "SITE001",
|
||||
"site_type": "office",
|
||||
"site_name": "Bureau Paris La Défense",
|
||||
"location": "Paris, France",
|
||||
"capacity_kw": 200,
|
||||
"status": "active",
|
||||
}
|
||||
|
||||
|
||||
def make_reading() -> dict[str, Any]:
|
||||
return {
|
||||
"timestamp": "2024-06-15T12:00:00Z",
|
||||
"site_id": "SITE001",
|
||||
"site_type": "office",
|
||||
"consumption_kw": 87.34,
|
||||
"consumption_kwh": 87.34,
|
||||
"voltage_v": 401.2,
|
||||
"current_a": 132.5,
|
||||
"power_factor": 0.923,
|
||||
"temperature_celsius": 22.1,
|
||||
"humidity_percent": 58.4,
|
||||
"null_reasons": [],
|
||||
"data_quality": "good",
|
||||
}
|
||||
|
||||
|
||||
async def test_fetch_sites_returns_sites() -> None:
|
||||
def handler(request: Request) -> Response:
|
||||
assert request.url.path == "/api/v1/sites"
|
||||
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_site()],
|
||||
)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
async with AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
) as client:
|
||||
sites = await fetch_sites(client)
|
||||
|
||||
assert len(sites) == 1
|
||||
assert sites[0]["site_id"] == "SITE001"
|
||||
assert sites[0]["site_type"] == "office"
|
||||
|
||||
|
||||
async def test_fetch_readings_sends_expected_query_parameters() -> None:
|
||||
captured_params: dict[str, str] = {}
|
||||
|
||||
def handler(request: Request) -> Response:
|
||||
nonlocal captured_params
|
||||
|
||||
captured_params = dict(request.url.params)
|
||||
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_reading()],
|
||||
)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
start_time = datetime.fromisoformat("2024-06-15T12:00:00")
|
||||
end_time = datetime.fromisoformat("2024-06-15T13:00:00")
|
||||
|
||||
async with AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
) as client:
|
||||
readings = await fetch_readings(
|
||||
client=client,
|
||||
site_id="SITE001",
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
limit=60,
|
||||
)
|
||||
|
||||
assert len(readings) == 1
|
||||
assert captured_params["site_id"] == "SITE001"
|
||||
assert captured_params["start_time"] == "2024-06-15T12:00:00"
|
||||
assert captured_params["end_time"] == "2024-06-15T13:00:00"
|
||||
assert captured_params["limit"] == "60"
|
||||
|
||||
|
||||
async def test_fetch_readings_rejects_non_list_response() -> None:
|
||||
def handler(request: Request) -> Response:
|
||||
return Response(
|
||||
status_code=200,
|
||||
json={"unexpected": "payload"},
|
||||
)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
async with AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
) as client:
|
||||
with pytest.raises(
|
||||
ValueError,
|
||||
match="La réponse /api/v1/readings doit être une liste",
|
||||
):
|
||||
await fetch_readings(
|
||||
client=client,
|
||||
site_id="SITE001",
|
||||
start_time=datetime.fromisoformat("2024-06-15T12:00:00"),
|
||||
end_time=datetime.fromisoformat("2024-06-15T13:00:00"),
|
||||
limit=60,
|
||||
)
|
||||
|
||||
|
||||
async def test_fetch_readings_raises_on_http_error() -> None:
|
||||
def handler(request: Request) -> Response:
|
||||
return Response(
|
||||
status_code=404,
|
||||
json={"detail": "Site non trouvé"},
|
||||
)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
async with AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
) as client:
|
||||
with pytest.raises(httpx.HTTPStatusError):
|
||||
await fetch_readings(
|
||||
client=client,
|
||||
site_id="SITE999",
|
||||
start_time=datetime.fromisoformat("2024-06-15T12:00:00"),
|
||||
end_time=datetime.fromisoformat("2024-06-15T13:00:00"),
|
||||
limit=60,
|
||||
)
|
||||
|
||||
|
||||
def test_build_reading_row_respects_database_contract() -> None:
|
||||
reading = make_reading()
|
||||
|
||||
row = build_reading_row(reading)
|
||||
|
||||
assert row["site_id"] == "SITE001"
|
||||
assert row["source"] == SOURCE_HISTORY
|
||||
assert row["source"] == "api_history"
|
||||
assert row["dataset_id"] is None
|
||||
|
||||
assert row["timestamp"] == datetime.fromisoformat("2024-06-15T12:00:00+00:00")
|
||||
|
||||
assert row["consumption_kw"] == 87.34
|
||||
assert row["consumption_kwh"] == 87.34
|
||||
assert row["data_quality"] == "good"
|
||||
assert row["null_reasons"] == []
|
||||
|
||||
assert row["imputed_values"] is None
|
||||
assert row["imputation_method"] is None
|
||||
|
||||
|
||||
def test_build_reading_row_keeps_null_values_and_quality() -> None:
|
||||
reading = make_reading()
|
||||
|
||||
reading["consumption_kw"] = None
|
||||
reading["consumption_kwh"] = None
|
||||
reading["voltage_v"] = None
|
||||
reading["current_a"] = None
|
||||
reading["power_factor"] = None
|
||||
reading["data_quality"] = "degraded"
|
||||
reading["null_reasons"] = [
|
||||
"consumption_sensor_failure",
|
||||
"electrical_sensor_failure",
|
||||
]
|
||||
|
||||
row = build_reading_row(reading)
|
||||
|
||||
assert row["consumption_kw"] is None
|
||||
assert row["consumption_kwh"] is None
|
||||
assert row["voltage_v"] is None
|
||||
assert row["current_a"] is None
|
||||
assert row["power_factor"] is None
|
||||
|
||||
assert row["data_quality"] == "degraded"
|
||||
assert row["null_reasons"] == [
|
||||
"consumption_sensor_failure",
|
||||
"electrical_sensor_failure",
|
||||
]
|
||||
|
||||
assert row["imputed_values"] is None
|
||||
assert row["imputation_method"] is None
|
||||
|
||||
|
||||
def test_build_reading_row_keeps_raw_source_data() -> None:
|
||||
reading = make_reading()
|
||||
|
||||
row = build_reading_row(reading)
|
||||
|
||||
raw_data = json.loads(row["raw_data"])
|
||||
|
||||
assert raw_data == reading
|
||||
|
||||
|
||||
def test_build_reading_batch_transforms_all_readings() -> None:
|
||||
first = make_reading()
|
||||
|
||||
second = make_reading()
|
||||
second["timestamp"] = "2024-06-15T12:01:00Z"
|
||||
second["consumption_kw"] = 90.5
|
||||
|
||||
rows = build_reading_batch([first, second])
|
||||
|
||||
assert len(rows) == 2
|
||||
|
||||
assert rows[0]["site_id"] == "SITE001"
|
||||
assert rows[0]["consumption_kw"] == 87.34
|
||||
|
||||
assert rows[1]["site_id"] == "SITE001"
|
||||
assert rows[1]["consumption_kw"] == 90.5
|
||||
|
||||
|
||||
def test_create_mock_api_client_requires_credentials(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
settings = SimpleNamespace(
|
||||
mock_api_username=None,
|
||||
mock_api_password=None,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"get_settings",
|
||||
lambda: settings,
|
||||
)
|
||||
|
||||
with pytest.raises(
|
||||
ValueError,
|
||||
match="Les identifiants de l'API Mock ne sont pas configurés",
|
||||
):
|
||||
mock_api_import.create_mock_api_client()
|
||||
|
||||
|
||||
async def test_create_mock_api_client_uses_configuration(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
password = MagicMock()
|
||||
password.get_secret_value.return_value = "test-password"
|
||||
|
||||
settings = SimpleNamespace(
|
||||
mock_api_base_url="https://mock.test/",
|
||||
mock_api_username="test-user",
|
||||
mock_api_password=password,
|
||||
mock_api_timeout_seconds=10.0,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"get_settings",
|
||||
lambda: settings,
|
||||
)
|
||||
|
||||
client = mock_api_import.create_mock_api_client()
|
||||
|
||||
try:
|
||||
assert str(client.base_url) == "https://mock.test"
|
||||
assert client.timeout.connect == 10.0
|
||||
finally:
|
||||
await client.aclose()
|
||||
|
||||
|
||||
async def test_upsert_sites_with_empty_list_does_nothing() -> None:
|
||||
connection = AsyncMock()
|
||||
|
||||
await upsert_sites(
|
||||
connection,
|
||||
[],
|
||||
)
|
||||
|
||||
connection.execute.assert_not_awaited()
|
||||
|
||||
|
||||
async def test_import_mock_api_history_dry_run_does_not_write(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
def handler(request: Request) -> Response:
|
||||
if request.url.path == "/api/v1/sites":
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_site()],
|
||||
)
|
||||
|
||||
if request.url.path == "/api/v1/readings":
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_reading()],
|
||||
)
|
||||
|
||||
return Response(status_code=404)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
client = AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"create_mock_api_client",
|
||||
lambda: client,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"get_settings",
|
||||
lambda: SimpleNamespace(
|
||||
database_url="postgresql+asyncpg://unused",
|
||||
),
|
||||
)
|
||||
|
||||
create_engine_mock = MagicMock()
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"create_async_engine",
|
||||
create_engine_mock,
|
||||
)
|
||||
|
||||
await mock_api_import.import_mock_api_history(
|
||||
start_time=datetime.fromisoformat("2024-06-15T12:00:00"),
|
||||
end_time=datetime.fromisoformat("2024-06-15T13:00:00"),
|
||||
limit=60,
|
||||
dry_run=True,
|
||||
)
|
||||
|
||||
create_engine_mock.assert_not_called()
|
||||
|
||||
|
||||
async def test_import_mock_api_history_loads_data(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
def handler(request: Request) -> Response:
|
||||
if request.url.path == "/api/v1/sites":
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_site()],
|
||||
)
|
||||
|
||||
if request.url.path == "/api/v1/readings":
|
||||
return Response(
|
||||
status_code=200,
|
||||
json=[make_reading()],
|
||||
)
|
||||
|
||||
return Response(status_code=404)
|
||||
|
||||
transport = MockTransport(handler)
|
||||
|
||||
client = AsyncClient(
|
||||
transport=transport,
|
||||
base_url="https://mock.test",
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"create_mock_api_client",
|
||||
lambda: client,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"get_settings",
|
||||
lambda: SimpleNamespace(
|
||||
database_url="postgresql+asyncpg://test:test@localhost/test",
|
||||
),
|
||||
)
|
||||
|
||||
connection = AsyncMock()
|
||||
|
||||
transaction_context = MagicMock()
|
||||
transaction_context.__aenter__ = AsyncMock(
|
||||
return_value=connection,
|
||||
)
|
||||
transaction_context.__aexit__ = AsyncMock(
|
||||
return_value=None,
|
||||
)
|
||||
|
||||
engine = MagicMock()
|
||||
engine.begin.return_value = transaction_context
|
||||
engine.dispose = AsyncMock()
|
||||
|
||||
create_engine_mock = MagicMock(
|
||||
return_value=engine,
|
||||
)
|
||||
|
||||
upsert_sites_mock = AsyncMock()
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"create_async_engine",
|
||||
create_engine_mock,
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"upsert_sites",
|
||||
upsert_sites_mock,
|
||||
)
|
||||
|
||||
await mock_api_import.import_mock_api_history(
|
||||
start_time=datetime.fromisoformat("2024-06-15T12:00:00"),
|
||||
end_time=datetime.fromisoformat("2024-06-15T13:00:00"),
|
||||
limit=60,
|
||||
dry_run=False,
|
||||
)
|
||||
|
||||
create_engine_mock.assert_called_once_with(
|
||||
"postgresql+asyncpg://test:test@localhost/test",
|
||||
pool_pre_ping=True,
|
||||
)
|
||||
|
||||
upsert_sites_mock.assert_awaited_once_with(
|
||||
connection,
|
||||
[make_site()],
|
||||
)
|
||||
|
||||
connection.execute.assert_awaited_once()
|
||||
engine.dispose.assert_awaited_once()
|
||||
|
||||
|
||||
def test_parse_datetime_accepts_z_suffix() -> None:
|
||||
result = mock_api_import.parse_datetime(
|
||||
"2024-06-15T12:00:00Z",
|
||||
)
|
||||
|
||||
assert result == datetime.fromisoformat(
|
||||
"2024-06-15T12:00:00+00:00",
|
||||
)
|
||||
|
||||
|
||||
def test_parse_args_reads_cli_parameters(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
monkeypatch.setattr(
|
||||
sys,
|
||||
"argv",
|
||||
[
|
||||
"mock_api_import",
|
||||
"--start-time",
|
||||
"2024-06-15T12:00:00Z",
|
||||
"--end-time",
|
||||
"2024-06-15T13:00:00Z",
|
||||
"--limit",
|
||||
"60",
|
||||
"--dry-run",
|
||||
],
|
||||
)
|
||||
|
||||
args = mock_api_import.parse_args()
|
||||
|
||||
assert args.start_time == datetime.fromisoformat(
|
||||
"2024-06-15T12:00:00+00:00",
|
||||
)
|
||||
assert args.end_time == datetime.fromisoformat(
|
||||
"2024-06-15T13:00:00+00:00",
|
||||
)
|
||||
assert args.limit == 60
|
||||
assert args.dry_run is True
|
||||
|
||||
|
||||
def test_main_rejects_limit_out_of_bounds(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
monkeypatch.setattr(
|
||||
sys,
|
||||
"argv",
|
||||
[
|
||||
"mock_api_import",
|
||||
"--start-time",
|
||||
"2024-06-15T12:00:00Z",
|
||||
"--end-time",
|
||||
"2024-06-15T13:00:00Z",
|
||||
"--limit",
|
||||
"0",
|
||||
],
|
||||
)
|
||||
|
||||
with pytest.raises(
|
||||
ValueError,
|
||||
match="--limit doit être compris entre 1 et 1000",
|
||||
):
|
||||
mock_api_import.main()
|
||||
|
||||
|
||||
def test_main_rejects_invalid_period(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
monkeypatch.setattr(
|
||||
sys,
|
||||
"argv",
|
||||
[
|
||||
"mock_api_import",
|
||||
"--start-time",
|
||||
"2024-06-15T14:00:00Z",
|
||||
"--end-time",
|
||||
"2024-06-15T13:00:00Z",
|
||||
"--limit",
|
||||
"60",
|
||||
],
|
||||
)
|
||||
|
||||
with pytest.raises(
|
||||
ValueError,
|
||||
match="--start-time doit être antérieur à --end-time",
|
||||
):
|
||||
mock_api_import.main()
|
||||
|
||||
|
||||
def test_main_runs_import(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
start_time = datetime.fromisoformat(
|
||||
"2024-06-15T12:00:00+00:00",
|
||||
)
|
||||
end_time = datetime.fromisoformat(
|
||||
"2024-06-15T13:00:00+00:00",
|
||||
)
|
||||
|
||||
import_mock = AsyncMock()
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"parse_args",
|
||||
lambda: SimpleNamespace(
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
limit=60,
|
||||
dry_run=True,
|
||||
),
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
mock_api_import,
|
||||
"import_mock_api_history",
|
||||
import_mock,
|
||||
)
|
||||
|
||||
mock_api_import.main()
|
||||
|
||||
import_mock.assert_awaited_once_with(
|
||||
start_time=start_time,
|
||||
end_time=end_time,
|
||||
limit=60,
|
||||
dry_run=True,
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.integration
|
||||
async def test_reading_insert_is_idempotent(
|
||||
session: AsyncSession,
|
||||
) -> None:
|
||||
reading = make_reading()
|
||||
row = build_reading_row(reading)
|
||||
|
||||
connection = await session.connection()
|
||||
|
||||
await upsert_sites(
|
||||
connection,
|
||||
[make_site()],
|
||||
)
|
||||
|
||||
await session.execute(
|
||||
READING_INSERT,
|
||||
[row],
|
||||
)
|
||||
|
||||
await session.execute(
|
||||
READING_INSERT,
|
||||
[row],
|
||||
)
|
||||
|
||||
result = await session.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT COUNT(*)
|
||||
FROM reading
|
||||
WHERE site_id = :site_id
|
||||
AND timestamp = :timestamp
|
||||
AND source = :source
|
||||
"""
|
||||
),
|
||||
{
|
||||
"site_id": row["site_id"],
|
||||
"timestamp": row["timestamp"],
|
||||
"source": row["source"],
|
||||
},
|
||||
)
|
||||
|
||||
assert result.scalar_one() == 1
|
||||
|
||||
await session.rollback()
|
||||
Generated
+2
-2
@@ -326,6 +326,7 @@ dependencies = [
|
||||
{ name = "argon2-cffi" },
|
||||
{ name = "asyncpg" },
|
||||
{ name = "fastapi" },
|
||||
{ name = "httpx" },
|
||||
{ name = "pandas" },
|
||||
{ name = "prometheus-fastapi-instrumentator" },
|
||||
{ name = "pydantic", extra = ["email"] },
|
||||
@@ -338,7 +339,6 @@ dependencies = [
|
||||
|
||||
[package.dev-dependencies]
|
||||
dev = [
|
||||
{ name = "httpx" },
|
||||
{ name = "mypy" },
|
||||
{ name = "pandas-stubs" },
|
||||
{ name = "pytest" },
|
||||
@@ -355,6 +355,7 @@ requires-dist = [
|
||||
{ name = "argon2-cffi", specifier = ">=23.1" },
|
||||
{ name = "asyncpg", specifier = ">=0.31.0" },
|
||||
{ name = "fastapi", specifier = ">=0.141.1" },
|
||||
{ name = "httpx", specifier = ">=0.28.1" },
|
||||
{ name = "pandas", specifier = ">=3.0.5" },
|
||||
{ name = "prometheus-fastapi-instrumentator", specifier = ">=8.1.0" },
|
||||
{ name = "pydantic", extras = ["email"], specifier = ">=2.13.5" },
|
||||
@@ -367,7 +368,6 @@ requires-dist = [
|
||||
|
||||
[package.metadata.requires-dev]
|
||||
dev = [
|
||||
{ name = "httpx", specifier = ">=0.28.1" },
|
||||
{ name = "mypy", specifier = ">=2.3.1" },
|
||||
{ name = "pandas-stubs", specifier = ">=3.0.5.260914" },
|
||||
{ name = "pytest", specifier = ">=9.1.1" },
|
||||
|
||||
@@ -50,6 +50,12 @@ services:
|
||||
APP_SECRET_KEY: ${APP_SECRET_KEY:?}
|
||||
APP_CORS_ORIGINS: ${APP_CORS_ORIGINS:-http://localhost:4200}
|
||||
DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
|
||||
|
||||
APP_MOCK_API_BASE_URL: ${APP_MOCK_API_BASE_URL:-https://api-mock.charlieandre.fr}
|
||||
APP_MOCK_API_USERNAME: ${APP_MOCK_API_USERNAME:-}
|
||||
APP_MOCK_API_PASSWORD: ${APP_MOCK_API_PASSWORD:-}
|
||||
APP_MOCK_API_TIMEOUT_SECONDS: ${APP_MOCK_API_TIMEOUT_SECONDS:-10}
|
||||
|
||||
APP_FRONTEND_RESET_PASSWORD_URL: ${APP_FRONTEND_RESET_PASSWORD_URL:-http://localhost:4200/reset-password}
|
||||
APP_SMTP_HOST: mailpit
|
||||
APP_SMTP_PORT: "1025"
|
||||
|
||||
+332
-69
@@ -6,10 +6,14 @@ système qui en découle.
|
||||
|
||||
## Ce que couvre ce document
|
||||
|
||||
**Dix tables applicatives existent** : quatre pour l'authentification, six pour les données
|
||||
d'énergie, dont l'hypertable `reading`. Les sections marquées `Fait` relèvent le code. Celles
|
||||
marquées `Cible` décrivent ce qui n'est pas écrit, au premier rang desquelles la chaîne
|
||||
d'ingestion, les agrégats continus, la compression et la rétention.
|
||||
**Douze tables applicatives existent** : six pour l'authentification et six pour les données
|
||||
d'énergie, dont l'hypertable `reading`.
|
||||
|
||||
Les sections marquées `Fait` relèvent du code déjà implémenté. Les sections marquées `Cible`
|
||||
décrivent les éléments prévus mais pas encore réalisés.
|
||||
|
||||
L'ingestion des deux sources de données du MVP est maintenant implémentée. L'orchestration
|
||||
Airflow, les agrégats continus, la compression et la rétention restent des cibles.
|
||||
|
||||
## Trois emplacements, trois rôles
|
||||
|
||||
@@ -35,8 +39,16 @@ Statut : `Fait`.
|
||||
- `db/init/100-extensions.sql` crée l'extension `timescaledb`.
|
||||
- `db/init/110-test-database.sql` crée `enervision_test`, dont le nom est attendu en dur par
|
||||
`apps/backend/tests/conftest.py`.
|
||||
- Cinq révisions Alembic. La première, `5353c0e4f094`, **ne crée aucune table** : elle
|
||||
établit `alembic_version` et refuse de s'appliquer si l'extension manque :
|
||||
- Six révisions Alembic sont actuellement appliquées.
|
||||
- La première, `5353c0e4f094`, **ne crée aucune table** : elle établit `alembic_version`
|
||||
et refuse de s'appliquer si l'extension TimescaleDB manque.
|
||||
- Les révisions suivantes créent les tables liées à l'authentification :
|
||||
`app_user`, `login_attempt`, `audit_log` et `refresh_token`.
|
||||
- La révision `e6d2026091501` crée les six tables Data et déclare l'hypertable `reading`.
|
||||
- La révision `c0adab96238c` ajoute les tables `password_reset_attempt`
|
||||
et `password_reset_token`.
|
||||
|
||||
La garde de la première migration est :
|
||||
|
||||
```sql
|
||||
IF NOT EXISTS (SELECT 1 FROM pg_extension WHERE extname = 'timescaledb') THEN
|
||||
@@ -47,37 +59,58 @@ END IF;
|
||||
Cette garde forme paire avec le 503 de `/api/v1/health/ready`. Un bootstrap sauté ne se voit pas
|
||||
au démarrage de l'API : ces deux gardes le rendent visible tôt, des deux côtés.
|
||||
|
||||
Les trois suivantes créent les tables de l'authentification, décrites plus bas : `app_user`,
|
||||
puis `login_attempt` et `audit_log`, puis `refresh_token`. La cinquième, `e6d2026091501`, crée
|
||||
les six tables de données décrites en fin de document et déclare l'hypertable `reading`.
|
||||
|
||||
## Cycle de vie d'une mesure
|
||||
|
||||
Statut : `Cible`, sauf l'hypertable `reading` qui existe. Ni l'ingestion, ni les agrégats
|
||||
continus, ni la compression, ni la rétention ne sont écrits.
|
||||
Statut : `Partiellement fait`.
|
||||
|
||||
Les mécanismes d'ingestion sont maintenant implémentés pour les deux sources de données du MVP :
|
||||
|
||||
- le dataset historique CSV/JSON avec `historical_import.py` ;
|
||||
- l'API Mock avec `mock_api_import.py`.
|
||||
|
||||
Les traitements sont actuellement exécutables directement depuis le backend.
|
||||
|
||||
L'orchestration avec Apache Airflow reste une cible, tout comme les agrégats continus,
|
||||
la compression et les politiques de rétention.
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
src["Source de mesures"] -.-> ing["Ingestion Airflow"]
|
||||
ing -.-> hy[("Hypertable reading")]
|
||||
csv["CSV + JSON"] --> hist["historical_import.py"]
|
||||
mock["API Mock"] --> api["mock_api_import.py"]
|
||||
|
||||
hist --> hy[("Hypertable reading")]
|
||||
api --> hy
|
||||
|
||||
airflow["Airflow"] -.-> hist
|
||||
airflow -.-> api
|
||||
|
||||
hy -.-> agg[("Agrégat continu")]
|
||||
hy -.-> comp["Compression"]
|
||||
hy -.-> ret["Rétention"]
|
||||
agg -.-> api["API FastAPI"]
|
||||
|
||||
agg -.-> backend["API FastAPI"]
|
||||
agg -.-> graf["Grafana"]
|
||||
```
|
||||
|
||||
Les lectures de l'API et de Grafana visent l'agrégat continu, pas la table brute : c'est tout
|
||||
l'intérêt de TimescaleDB, et cela doit rester vrai quand les volumes augmenteront.
|
||||
Les flèches pleines représentent les traitements actuellement implémentés.
|
||||
|
||||
Les flèches pointillées représentent les éléments encore prévus comme cibles.
|
||||
|
||||
Les lectures futures de l'API et de Grafana visent l'agrégat continu plutôt que la table brute
|
||||
lorsque cette partie TimescaleDB sera mise en place.
|
||||
|
||||
## Tables d'authentification
|
||||
|
||||
Statut : `Fait`. Elles ne sont pas des séries temporelles et n'ont donc rien à voir avec les
|
||||
hypertables ; elles vivent dans `apps/backend/alembic/`, qui porte le schéma exposé par l'API.
|
||||
Statut : `Fait`.
|
||||
|
||||
Elles ne sont pas des séries temporelles et n'ont donc rien à voir avec les hypertables ;
|
||||
elles vivent dans `apps/backend/alembic/`, qui porte le schéma exposé par l'API.
|
||||
|
||||
```mermaid
|
||||
erDiagram
|
||||
APP_USER ||--o{ REFRESH_TOKEN : ouvre
|
||||
APP_USER ||--o{ PASSWORD_RESET_TOKEN : recoit
|
||||
|
||||
APP_USER {
|
||||
uuid id PK
|
||||
string email UK
|
||||
@@ -88,6 +121,7 @@ erDiagram
|
||||
bool must_change_password
|
||||
timestamptz credentials_changed_at
|
||||
}
|
||||
|
||||
REFRESH_TOKEN {
|
||||
uuid id PK
|
||||
uuid family_id
|
||||
@@ -99,6 +133,7 @@ erDiagram
|
||||
text revoked_reason
|
||||
uuid replaced_by
|
||||
}
|
||||
|
||||
LOGIN_ATTEMPT {
|
||||
bigint id PK
|
||||
timestamptz occurred_at
|
||||
@@ -106,6 +141,7 @@ erDiagram
|
||||
inet client_ip
|
||||
text outcome
|
||||
}
|
||||
|
||||
AUDIT_LOG {
|
||||
bigint id PK
|
||||
timestamptz occurred_at
|
||||
@@ -114,31 +150,47 @@ erDiagram
|
||||
text action
|
||||
jsonb detail
|
||||
}
|
||||
|
||||
PASSWORD_RESET_ATTEMPT {
|
||||
bigint id PK
|
||||
timestamptz occurred_at
|
||||
string email_tried
|
||||
inet client_ip
|
||||
}
|
||||
|
||||
PASSWORD_RESET_TOKEN {
|
||||
uuid id PK
|
||||
uuid user_id FK
|
||||
bytea token_hash UK
|
||||
timestamptz issued_at
|
||||
timestamptz expires_at
|
||||
timestamptz consumed_at
|
||||
inet client_ip
|
||||
text user_agent
|
||||
}
|
||||
```
|
||||
|
||||
Quatre choix de modélisation portent une intention et se défendent seuls :
|
||||
Plusieurs choix de modélisation portent une intention précise :
|
||||
|
||||
- **`app_user` et non `user`** : `user` est un mot réservé PostgreSQL, raccourci de
|
||||
`CURRENT_USER`. Le nom rappelle en prime qu'il s'agit d'un compte applicatif, par opposition
|
||||
au rôle PostgreSQL qui portera le cantonnement de l'ETL.
|
||||
- **`credentials_changed_at`, une seule colonne**, couvre le changement de mot de passe, le
|
||||
changement de rôle et la désactivation. Un compteur de version ne dirait rien à un humain qui
|
||||
lit un audit.
|
||||
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation. S'il
|
||||
glissait, la promesse de sept jours serait fictive et une session active ne finirait jamais.
|
||||
- **`audit_log.actor_id` n'a aucune clé étrangère**, et `actor_email` comme `actor_role` sont
|
||||
dénormalisés. Une contrainte `ON DELETE SET NULL` déclencherait un `UPDATE` que le déclencheur
|
||||
d'ajout seul refuserait. Voir l'[ADR 0004](../adr/0004-journal-d-audit-en-ajout-seul.md).
|
||||
`CURRENT_USER`. Le nom rappelle aussi qu'il s'agit d'un compte applicatif.
|
||||
- **`credentials_changed_at`, une seule colonne**, couvre notamment le changement de mot de passe,
|
||||
le changement de rôle et la désactivation.
|
||||
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation.
|
||||
- **`audit_log.actor_id` n'a aucune clé étrangère** afin de conserver les informations d'audit
|
||||
même si l'entité d'origine évolue.
|
||||
- `password_reset_token` ne stocke que l'empreinte du jeton et jamais sa valeur directement.
|
||||
- `password_reset_attempt` est séparée de `audit_log`, car son volume peut être piloté
|
||||
par des demandes externes répétées.
|
||||
|
||||
`audit_log` porte deux déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`. Elle n'est
|
||||
donc **pas** une hypertable : une politique de rétention émettrait des `DELETE` qu'ils
|
||||
refuseraient. `login_attempt`, à l'inverse, est faite pour se purger, puisque son volume est
|
||||
piloté par l'attaquant.
|
||||
`audit_log` porte des déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`.
|
||||
Elle n'est donc **pas** une hypertable.
|
||||
|
||||
## Gabarit de révision créant une hypertable
|
||||
|
||||
Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision. La révision
|
||||
`e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
||||
Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision.
|
||||
|
||||
La révision `e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
||||
|
||||
```python
|
||||
def upgrade() -> None:
|
||||
@@ -182,24 +234,28 @@ colonne de temps : les index déclarés dans la révision le couvrent déjà.
|
||||
|
||||
## Questions ouvertes
|
||||
|
||||
Elles relèvent du jalon J2, « valider le périmètre retenu ». Le schéma est livré : ce qui suit
|
||||
porte sur son exploitation, plus sur sa forme.
|
||||
Elles portent maintenant principalement sur l'exploitation du schéma :
|
||||
|
||||
- **Quelle granularité** à l'ingestion : la seconde, la minute, le quart d'heure.
|
||||
- **Quels agrégats continus**, et sur quelles fenêtres.
|
||||
- **Quelle profondeur de rétention** en données brutes, et à partir de quand on compresse.
|
||||
- **Multi-tenant ou non** : un site appartient-il à un client, et faut-il cloisonner les lectures.
|
||||
- **Quelle granularité** conserver à long terme à l'ingestion : seconde, minute ou quart d'heure.
|
||||
- **Quels agrégats continus** créer et sur quelles fenêtres.
|
||||
- **Quelle profondeur de rétention** conserver en données brutes et à partir de quand compresser.
|
||||
- **Multi-tenant ou non** : un site appartient-il à un client et faut-il cloisonner les lectures.
|
||||
|
||||
## Modélisation détaillée des données
|
||||
|
||||
Cette modélisation prend en compte les fichiers CSV historiques,
|
||||
leurs métadonnées JSON et les données de l’API Mock.
|
||||
Elle comprend six tables, depuis le stockage des mesures
|
||||
jusqu’aux recommandations proposées à l’utilisateur.
|
||||
Cette modélisation prend en compte :
|
||||
|
||||
- les fichiers CSV historiques ;
|
||||
- leurs métadonnées JSON ;
|
||||
- les données de l'API Mock.
|
||||
|
||||
Elle comprend six tables Data, depuis le stockage des mesures jusqu'aux recommandations proposées
|
||||
à l'utilisateur.
|
||||
|
||||
### Schéma de données
|
||||
|
||||
Le diagramme ci-dessous présente les tables et leurs relations.
|
||||
|
||||
La révision `e6d2026091501` les crée.
|
||||
|
||||

|
||||
@@ -208,21 +264,20 @@ La révision `e6d2026091501` les crée.
|
||||
|
||||
### Description des tables
|
||||
|
||||
Chaque table remplit un rôle précis dans le traitement et l’exploitation
|
||||
des données.
|
||||
Chaque table remplit un rôle précis dans le traitement et l'exploitation des données.
|
||||
|
||||
| Table | Rôle | Origine des informations |
|
||||
|---|---|---|
|
||||
| `dataset` | Identifier les jeux historiques, retrouver leurs fichiers et conserver leurs métadonnées | Archive CSV/JSON et informations ajoutées lors de l’import |
|
||||
| `dataset` | Identifier les jeux historiques, retrouver leurs fichiers et conserver leurs métadonnées | Archive CSV/JSON et informations ajoutées lors de l'import |
|
||||
| `site` | Regrouper les informations des sites : identifiant, nom, type et caractéristiques disponibles | CSV et API Mock `/api/v1/sites` |
|
||||
| `reading` | Stocker les mesures, leur provenance, leur qualité et les éventuelles valeurs imputées | CSV et API Mock `/current` et `/readings` |
|
||||
| `prediction` | Conserver les prévisions, leur période cible et la référence du modèle utilisé | Traitements ML d’EnerVision |
|
||||
| `prediction` | Conserver les prévisions, leur période cible et la référence du modèle utilisé | Traitements ML d'EnerVision |
|
||||
| `alert` | Enregistrer les alertes, leur type, leur gravité et leur message | API Mock `/alerts` et détections EnerVision |
|
||||
| `recommendation` | Proposer des actions et expliquer la règle qui les motive | Règles métier d’EnerVision |
|
||||
| `recommendation` | Proposer des actions et expliquer la règle qui les motive | Règles métier d'EnerVision |
|
||||
|
||||
Les anomalies historiques décrites dans les JSON sont conservées
|
||||
dans `dataset.metadata`. Elles servent à l’analyse des données
|
||||
et ne sont pas considérées comme des alertes actuelles.
|
||||
Les anomalies historiques décrites dans les JSON sont conservées dans `dataset.metadata`.
|
||||
|
||||
Elles servent à l'analyse des données et ne sont pas considérées comme des alertes actuelles.
|
||||
|
||||
### Relations entre les tables
|
||||
|
||||
@@ -232,15 +287,21 @@ et ne sont pas considérées comme des alertes actuelles.
|
||||
- Une alerte peut être associée à une prévision du même site.
|
||||
- Une alerte peut donner lieu à plusieurs recommandations.
|
||||
|
||||
## Ingestion des données historiques
|
||||
# Ingestion des données historiques
|
||||
|
||||
Le MVP EnerVision initialise les données énergétiques à partir du dataset fourni dans le cadre du projet.
|
||||
Statut : `Fait`.
|
||||
|
||||
Le dataset de référence contient 122 647 mesures issues de 7 sites et couvre la période du 1er janvier 2023 au 31 décembre 2024.
|
||||
Le MVP EnerVision initialise les données énergétiques à partir du dataset fourni dans le cadre
|
||||
du projet.
|
||||
|
||||
Les fichiers sources CSV et JSON sont nécessaires uniquement pour l'initialisation des données. Ils ne sont pas versionnés dans Git et sont placés localement dans `data/raw/`.
|
||||
Le dataset de référence contient 122 647 mesures issues de 7 sites et couvre la période
|
||||
du 1er janvier 2023 au 31 décembre 2024.
|
||||
|
||||
### Architecture du flux
|
||||
Les fichiers sources CSV et JSON sont nécessaires uniquement pour l'initialisation des données.
|
||||
|
||||
Ils ne sont pas versionnés dans Git et sont placés localement dans `data/raw/`.
|
||||
|
||||
## Architecture du flux historique
|
||||
|
||||
```text
|
||||
Dataset CSV + métadonnées JSON
|
||||
@@ -271,17 +332,26 @@ Dataset CSV + métadonnées JSON
|
||||
|
||||
Le pipeline est développé en Python.
|
||||
|
||||
Pandas est utilisé pour l'extraction, la validation et la préparation des données. SQLAlchemy Async assure le chargement transactionnel dans PostgreSQL/TimescaleDB.
|
||||
Pandas est utilisé pour l'extraction, la validation et la préparation des données.
|
||||
|
||||
SQLAlchemy Async assure le chargement transactionnel dans PostgreSQL/TimescaleDB.
|
||||
|
||||
Une empreinte SHA-256 permet d'identifier le dataset utilisé et d'assurer sa traçabilité.
|
||||
|
||||
Les valeurs manquantes sont conservées pendant l'ingestion afin de préserver les données sources. Aucune imputation n'est réalisée à cette étape.
|
||||
Les valeurs manquantes sont conservées pendant l'ingestion afin de préserver les données sources.
|
||||
|
||||
Aucune imputation n'est réalisée à cette étape.
|
||||
|
||||
Le chargement des mesures est effectué par batches de 1 000 lignes.
|
||||
|
||||
Les données provenant du dataset CSV sont identifiées par `source = "csv"` et associées à leur `dataset_id`.
|
||||
Les données provenant du dataset CSV sont identifiées par :
|
||||
|
||||
### Résultats validés
|
||||
```text
|
||||
source = "csv"
|
||||
dataset_id = identifiant du dataset
|
||||
```
|
||||
|
||||
## Résultats validés pour l'historique
|
||||
|
||||
Le chargement de référence a permis d'obtenir :
|
||||
|
||||
@@ -290,14 +360,207 @@ Le chargement de référence a permis d'obtenir :
|
||||
- 122 647 mesures ;
|
||||
- 0 doublon détecté dans le dataset source.
|
||||
|
||||
L'idempotence a également été vérifiée par une deuxième exécution du pipeline : aucune nouvelle mesure n'a été créée et le nombre de `reading` est resté à 122 647.
|
||||
L'idempotence a également été vérifiée par une deuxième exécution du pipeline :
|
||||
aucune nouvelle mesure n'a été créée et le nombre de `reading` est resté à 122 647.
|
||||
|
||||
La procédure détaillée d'installation, d'exécution, de validation et de contrôle du pipeline est disponible dans `etl/README.md`.
|
||||
La procédure détaillée d'installation, d'exécution, de validation et de contrôle du pipeline
|
||||
est disponible dans `etl/README.md`.
|
||||
|
||||
### Évolution prévue
|
||||
# Ingestion depuis l'API Mock
|
||||
|
||||
L'étape suivante consiste à orchestrer les traitements Data avec Apache Airflow.
|
||||
Statut : `Fait`.
|
||||
|
||||
L'orchestration réutilisera la logique ETL existante afin de séparer la logique de traitement de la planification, du suivi des exécutions et de la gestion des erreurs.
|
||||
La deuxième source du pipeline Data est l'API Mock EnerVision.
|
||||
|
||||
Le pipeline servira ensuite de base à la préparation des données nécessaires au modèle de Machine Learning.
|
||||
Le traitement est implémenté dans :
|
||||
|
||||
```text
|
||||
apps/backend/app/etl/mock_api_import.py
|
||||
```
|
||||
|
||||
## Endpoints utilisés
|
||||
|
||||
Le pipeline récupère les informations des sites depuis :
|
||||
|
||||
```text
|
||||
GET /api/v1/sites
|
||||
```
|
||||
|
||||
puis les mesures historiques simulées depuis :
|
||||
|
||||
```text
|
||||
GET /api/v1/readings
|
||||
```
|
||||
|
||||
Pour `/api/v1/readings`, les informations suivantes sont envoyées :
|
||||
|
||||
```text
|
||||
site_id
|
||||
start_time
|
||||
end_time
|
||||
limit
|
||||
```
|
||||
|
||||
Les paramètres de ligne de commande disponibles pour l'import sont :
|
||||
|
||||
```text
|
||||
--start-time
|
||||
--end-time
|
||||
--limit
|
||||
--dry-run
|
||||
```
|
||||
|
||||
## Flux d'ingestion API Mock
|
||||
|
||||
```text
|
||||
API Mock
|
||||
|
|
||||
+-----+------+
|
||||
| |
|
||||
v v
|
||||
/sites /readings
|
||||
| |
|
||||
+-----+------+
|
||||
|
|
||||
v
|
||||
mock_api_import.py
|
||||
|
|
||||
v
|
||||
Transformation
|
||||
+ qualité data
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
| |
|
||||
v v
|
||||
site reading
|
||||
```
|
||||
|
||||
Les informations des sites sont insérées ou mises à jour dans `site`.
|
||||
|
||||
Les mesures sont enregistrées dans l'hypertable `reading` avec :
|
||||
|
||||
```text
|
||||
source = "api_history"
|
||||
dataset_id = NULL
|
||||
```
|
||||
|
||||
Les données provenant de l'API Mock ne sont donc pas associées à un enregistrement de la table
|
||||
`dataset`.
|
||||
|
||||
La réponse source reçue depuis l'API est conservée dans :
|
||||
|
||||
```text
|
||||
raw_data
|
||||
```
|
||||
|
||||
## Qualité des données de l'API Mock
|
||||
|
||||
Les valeurs `NULL` ne sont pas remplacées pendant l'ingestion.
|
||||
|
||||
Les informations suivantes fournies par l'API sont conservées :
|
||||
|
||||
```text
|
||||
data_quality
|
||||
null_reasons
|
||||
```
|
||||
|
||||
Cette conservation permet de distinguer une valeur manquante d'une valeur réelle égale à zéro
|
||||
et de garder les informations liées aux éventuelles défaillances de capteurs.
|
||||
|
||||
Aucune imputation n'est réalisée pendant cette phase :
|
||||
|
||||
```text
|
||||
imputed_values = NULL
|
||||
imputation_method = NULL
|
||||
```
|
||||
|
||||
## Validation de l'import API Mock
|
||||
|
||||
Un scénario de validation a été exécuté pour les 7 sites sur la période :
|
||||
|
||||
```text
|
||||
15/06/2024 12:00 UTC
|
||||
à
|
||||
15/06/2024 13:00 UTC
|
||||
```
|
||||
|
||||
avec :
|
||||
|
||||
```text
|
||||
limit = 60
|
||||
```
|
||||
|
||||
Résultat :
|
||||
|
||||
```text
|
||||
7 sites
|
||||
60 lectures par site
|
||||
420 lectures récupérées
|
||||
```
|
||||
|
||||
Les données ont été chargées dans PostgreSQL/TimescaleDB puis contrôlées directement en base.
|
||||
|
||||
Les contrôles ont confirmé :
|
||||
|
||||
- `source = "api_history"` ;
|
||||
- `dataset_id = NULL` ;
|
||||
- la conservation des valeurs `NULL` ;
|
||||
- la conservation de `data_quality` ;
|
||||
- la conservation de `null_reasons` ;
|
||||
- la conservation de `raw_data`.
|
||||
|
||||
L'idempotence a été vérifiée en rejouant le même import.
|
||||
|
||||
Une mesure déjà présente n'est pas ajoutée une seconde fois.
|
||||
|
||||
Les tests automatisés couvrent également :
|
||||
|
||||
- la récupération des sites ;
|
||||
- les paramètres envoyés à `/api/v1/readings` ;
|
||||
- les réponses HTTP en erreur ;
|
||||
- le format de la réponse ;
|
||||
- la transformation des mesures ;
|
||||
- les valeurs manquantes ;
|
||||
- la qualité des données ;
|
||||
- la conservation des données sources ;
|
||||
- l'idempotence en base.
|
||||
|
||||
# Évolution prévue
|
||||
|
||||
La prochaine étape consiste à orchestrer les deux mécanismes d'ingestion avec Apache Airflow.
|
||||
|
||||
```text
|
||||
CSV / JSON ----------------+
|
||||
|
|
||||
v
|
||||
+------------------+
|
||||
| Airflow |
|
||||
+------------------+
|
||||
|
|
||||
+----------------+----------------+
|
||||
| |
|
||||
v v
|
||||
historical_import.py mock_api_import.py
|
||||
| |
|
||||
+----------------+----------------+
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
```
|
||||
|
||||
Airflow servira à :
|
||||
|
||||
- planifier les traitements ;
|
||||
- définir leur ordre d'exécution ;
|
||||
- suivre leur état ;
|
||||
- gérer et remonter les erreurs ;
|
||||
- faciliter les exécutions récurrentes.
|
||||
|
||||
Airflow ne remplacera pas la logique ETL déjà implémentée.
|
||||
|
||||
Les scripts Python resteront responsables de l'extraction, de la validation, de la transformation
|
||||
et du chargement des données.
|
||||
|
||||
Le pipeline servira ensuite de base à la préparation des données nécessaires au modèle
|
||||
de Machine Learning.
|
||||
+307
-22
@@ -2,9 +2,10 @@
|
||||
|
||||
## Objectif
|
||||
|
||||
Le pipeline ETL EnerVision permet d'intégrer les données énergétiques historiques dans PostgreSQL/TimescaleDB.
|
||||
Le pipeline ETL EnerVision permet d'intégrer les données énergétiques dans PostgreSQL/TimescaleDB à partir de deux sources :
|
||||
|
||||
Cette première étape du pipeline Data permet de charger le dataset fourni dans le cadre du projet, contenant les mesures énergétiques de 7 sites sur la période du 1er janvier 2023 au 31 décembre 2024.
|
||||
- le dataset historique CSV/JSON fourni dans le cadre du projet ;
|
||||
- l'API Mock EnerVision.
|
||||
|
||||
Le pipeline assure :
|
||||
|
||||
@@ -12,12 +13,15 @@ Le pipeline assure :
|
||||
- la validation de leur structure et de leur cohérence ;
|
||||
- la normalisation des données nécessaires au stockage ;
|
||||
- le suivi de la qualité des données ;
|
||||
- la traçabilité du dataset importé ;
|
||||
- la traçabilité des données importées ;
|
||||
- le chargement des données dans PostgreSQL/TimescaleDB ;
|
||||
- la conservation des valeurs manquantes et des informations de qualité ;
|
||||
- l'idempotence du chargement afin d'éviter la création de doublons.
|
||||
|
||||
## Données sources
|
||||
|
||||
### Dataset historique
|
||||
|
||||
Le dataset est fourni par le formateur dans le cadre du projet EnerVision.
|
||||
|
||||
Il contient les deux fichiers suivants :
|
||||
@@ -29,7 +33,7 @@ dataset_metadata.json
|
||||
|
||||
Ces fichiers sont nécessaires une seule fois pour initialiser les données historiques de l'environnement.
|
||||
|
||||
Ils ne sont pas versionnés dans Git. Chaque membre de l'équipe récupère manuellement une fois les fichiers fournis par le formateur et les place dans :
|
||||
Ils ne sont pas versionnés dans Git. Chaque membre de l'équipe récupère manuellement les fichiers fournis par le formateur et les place dans :
|
||||
|
||||
```text
|
||||
data/raw/
|
||||
@@ -47,14 +51,26 @@ data/
|
||||
|
||||
Le fichier `.gitkeep` est versionné afin de conserver le répertoire `data/raw/` dans Git. Les fichiers CSV et JSON sont ignorés par Git.
|
||||
|
||||
### API Mock
|
||||
|
||||
La deuxième source est l'API Mock EnerVision.
|
||||
|
||||
Elle permet de récupérer :
|
||||
|
||||
- les informations des sites avec `GET /api/v1/sites` ;
|
||||
- les mesures simulées avec `GET /api/v1/readings`.
|
||||
|
||||
L'API Mock est utilisée pour compléter les données historiques avec des mesures simulées récupérées sur une période donnée.
|
||||
|
||||
## Technologies utilisées
|
||||
|
||||
| Technologie | Utilisation |
|
||||
|---|---|
|
||||
| Python | Développement du pipeline ETL |
|
||||
| Pandas | Lecture, validation et transformation des données |
|
||||
| JSON | Lecture des métadonnées du dataset |
|
||||
| hashlib / SHA-256 | Identification, intégrité et traçabilité du dataset |
|
||||
| Pandas | Lecture, validation et transformation du dataset historique |
|
||||
| JSON | Lecture des métadonnées et conservation des données sources |
|
||||
| HTTPX | Appels HTTP asynchrones vers l'API Mock |
|
||||
| hashlib / SHA-256 | Identification, intégrité et traçabilité du dataset historique |
|
||||
| SQLAlchemy Async | Connexion et chargement asynchrone en base |
|
||||
| PostgreSQL | Stockage relationnel |
|
||||
| TimescaleDB | Stockage des séries temporelles énergétiques |
|
||||
@@ -62,11 +78,14 @@ Le fichier `.gitkeep` est versionné afin de conserver le répertoire `data/raw/
|
||||
| Alembic | Gestion des migrations du schéma |
|
||||
| uv | Gestion et exécution de l'environnement Python |
|
||||
| Ruff | Contrôle de la qualité du code |
|
||||
| mypy | Vérification du typage |
|
||||
| Pytest | Tests automatisés |
|
||||
|
||||
## Fonctionnement du pipeline
|
||||
# Import du dataset historique
|
||||
|
||||
Le script principal d'import se trouve dans :
|
||||
## Fonctionnement du pipeline historique
|
||||
|
||||
Le script d'import se trouve dans :
|
||||
|
||||
```text
|
||||
apps/backend/app/etl/historical_import.py
|
||||
@@ -207,7 +226,7 @@ Valeurs manquantes identifiées :
|
||||
| `humidity_percent` | 3 423 |
|
||||
| `solar_irradiance_wm2` | 3 964 |
|
||||
|
||||
## Exécution en dry-run
|
||||
## Exécution historique en dry-run
|
||||
|
||||
Depuis le dossier :
|
||||
|
||||
@@ -227,7 +246,7 @@ uv run python -m app.etl.historical_import `
|
||||
|
||||
Aucune donnée n'est écrite dans la base pendant cette exécution.
|
||||
|
||||
## Chargement réel
|
||||
## Chargement historique réel
|
||||
|
||||
Depuis `apps/backend/` :
|
||||
|
||||
@@ -249,7 +268,7 @@ Chargement : 2000/122647
|
||||
Chargement : 122647/122647
|
||||
```
|
||||
|
||||
## Résultats obtenus
|
||||
## Résultats obtenus pour le dataset historique
|
||||
|
||||
Après le chargement initial, les contrôles en base ont confirmé :
|
||||
|
||||
@@ -266,7 +285,7 @@ Le premier import a créé :
|
||||
nouvelles lectures : 122647
|
||||
```
|
||||
|
||||
## Idempotence
|
||||
## Idempotence du dataset historique
|
||||
|
||||
Le pipeline a été exécuté une deuxième fois avec exactement le même dataset afin de vérifier son idempotence.
|
||||
|
||||
@@ -280,7 +299,7 @@ nouvelles lectures : 0
|
||||
|
||||
Une nouvelle exécution du même import ne crée donc pas de mesures supplémentaires pour le dataset testé.
|
||||
|
||||
## Vérifications SQL
|
||||
## Vérifications SQL du dataset historique
|
||||
|
||||
Depuis la racine du projet, vérifier le nombre d'enregistrements avec :
|
||||
|
||||
@@ -302,21 +321,213 @@ Vérifier la source des mesures avec :
|
||||
docker compose exec db psql -U enervision -d enervision -c "SELECT source, COUNT(*) FROM reading GROUP BY source ORDER BY source;"
|
||||
```
|
||||
|
||||
Résultat attendu :
|
||||
Résultat attendu pour le dataset historique :
|
||||
|
||||
```text
|
||||
csv | 122647
|
||||
```
|
||||
|
||||
## Tests et qualité
|
||||
# Import depuis l'API Mock
|
||||
|
||||
Les tests automatisés du pipeline sont situés dans :
|
||||
## Fonctionnement
|
||||
|
||||
Le script d'import de l'API Mock se trouve dans :
|
||||
|
||||
```text
|
||||
apps/backend/app/etl/mock_api_import.py
|
||||
```
|
||||
|
||||
Le flux est le suivant :
|
||||
|
||||
```text
|
||||
API Mock
|
||||
|
|
||||
+-----+------+
|
||||
| |
|
||||
v v
|
||||
/sites /readings
|
||||
| |
|
||||
+-----+------+
|
||||
|
|
||||
v
|
||||
mock_api_import.py
|
||||
|
|
||||
v
|
||||
Transformation
|
||||
+ qualité data
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
| |
|
||||
v v
|
||||
site reading
|
||||
```
|
||||
|
||||
Le pipeline commence par récupérer les sites avec :
|
||||
|
||||
```text
|
||||
GET /api/v1/sites
|
||||
```
|
||||
|
||||
Il récupère ensuite les mesures de chaque site avec :
|
||||
|
||||
```text
|
||||
GET /api/v1/readings
|
||||
```
|
||||
|
||||
Les paramètres envoyés à `/api/v1/readings` sont :
|
||||
|
||||
```text
|
||||
site_id
|
||||
start_time
|
||||
end_time
|
||||
limit
|
||||
```
|
||||
|
||||
Le paramètre `limit` doit être compris entre 1 et 1000.
|
||||
|
||||
## Configuration de l'API Mock
|
||||
|
||||
La connexion à l'API Mock est configurée avec les variables d'environnement suivantes :
|
||||
|
||||
```text
|
||||
APP_MOCK_API_BASE_URL
|
||||
APP_MOCK_API_USERNAME
|
||||
APP_MOCK_API_PASSWORD
|
||||
APP_MOCK_API_TIMEOUT_SECONDS
|
||||
```
|
||||
|
||||
Les identifiants réels ne sont pas versionnés dans Git.
|
||||
|
||||
Les fichiers `.env.example` indiquent uniquement les variables nécessaires à l'exécution.
|
||||
|
||||
## Transformation des mesures API
|
||||
|
||||
Les mesures provenant de l'API Mock sont enregistrées dans `reading` avec :
|
||||
|
||||
```text
|
||||
source = "api_history"
|
||||
dataset_id = NULL
|
||||
```
|
||||
|
||||
Les mesures provenant de l'API ne sont donc pas rattachées à un dataset historique.
|
||||
|
||||
Le timestamp reçu depuis l'API est converti en `datetime` avec timezone avant le chargement.
|
||||
|
||||
La réponse source est conservée dans :
|
||||
|
||||
```text
|
||||
raw_data
|
||||
```
|
||||
|
||||
afin de préserver la donnée reçue et faciliter la traçabilité.
|
||||
|
||||
## Qualité des données API
|
||||
|
||||
Les valeurs `NULL` fournies par l'API sont conservées telles quelles.
|
||||
|
||||
Une valeur manquante n'est pas transformée en zéro et la mesure n'est pas supprimée.
|
||||
|
||||
Le pipeline conserve également :
|
||||
|
||||
```text
|
||||
data_quality
|
||||
null_reasons
|
||||
```
|
||||
|
||||
Les niveaux de qualité possibles sont :
|
||||
|
||||
```text
|
||||
good
|
||||
partial
|
||||
degraded
|
||||
critical
|
||||
```
|
||||
|
||||
Aucune imputation n'est réalisée pendant l'ingestion :
|
||||
|
||||
```text
|
||||
imputed_values = NULL
|
||||
imputation_method = NULL
|
||||
```
|
||||
|
||||
Cette stratégie permet de distinguer une véritable valeur nulle ou manquante d'une consommation égale à zéro et de conserver les informations liées aux défaillances de capteurs.
|
||||
|
||||
## Dry-run de l'API Mock
|
||||
|
||||
Le mode `--dry-run` permet de tester la connexion, la récupération des sites et la récupération des mesures sans écrire dans PostgreSQL.
|
||||
|
||||
Depuis `apps/backend/` :
|
||||
|
||||
```powershell
|
||||
uv run python -m app.etl.mock_api_import `
|
||||
--start-time "2024-06-15T12:00:00" `
|
||||
--end-time "2024-06-15T13:00:00" `
|
||||
--limit 60 `
|
||||
--dry-run
|
||||
```
|
||||
|
||||
## Chargement réel depuis l'API Mock
|
||||
|
||||
Depuis `apps/backend/` :
|
||||
|
||||
```powershell
|
||||
uv run python -m app.etl.mock_api_import `
|
||||
--start-time "2024-06-15T12:00:00" `
|
||||
--end-time "2024-06-15T13:00:00" `
|
||||
--limit 60
|
||||
```
|
||||
|
||||
## Résultat validé pour l'API Mock
|
||||
|
||||
Le scénario de validation utilisé couvre la période :
|
||||
|
||||
```text
|
||||
15/06/2024 12:00 UTC
|
||||
à
|
||||
15/06/2024 13:00 UTC
|
||||
```
|
||||
|
||||
avec une limite de 60 lectures par site.
|
||||
|
||||
Résultat obtenu :
|
||||
|
||||
```text
|
||||
sites récupérés : 7
|
||||
lectures par site : 60
|
||||
lectures récupérées : 420
|
||||
source : api_history
|
||||
dataset_id : NULL
|
||||
```
|
||||
|
||||
Les contrôles effectués directement dans PostgreSQL/TimescaleDB ont confirmé :
|
||||
|
||||
- l'enregistrement des mesures dans `reading` ;
|
||||
- la présence des 7 sites ;
|
||||
- `source = "api_history"` ;
|
||||
- `dataset_id = NULL` ;
|
||||
- la conservation des valeurs `NULL` ;
|
||||
- la conservation de `data_quality` ;
|
||||
- la conservation de `null_reasons` ;
|
||||
- la conservation de la donnée source dans `raw_data`.
|
||||
|
||||
## Idempotence de l'import API Mock
|
||||
|
||||
Le même import a été exécuté plusieurs fois afin de vérifier qu'une mesure déjà présente n'est pas créée une seconde fois.
|
||||
|
||||
L'idempotence repose sur la contrainte d'unicité de la table `reading` et sur la gestion des conflits lors de l'insertion.
|
||||
|
||||
Un test d'intégration automatisé vérifie également ce comportement.
|
||||
|
||||
# Tests et qualité
|
||||
|
||||
Les tests automatisés des pipelines ETL sont situés dans :
|
||||
|
||||
```text
|
||||
apps/backend/tests/etl/
|
||||
```
|
||||
|
||||
Ils couvrent notamment :
|
||||
Les tests de l'import historique couvrent notamment :
|
||||
|
||||
- la validation du dataset ;
|
||||
- les colonnes obligatoires ;
|
||||
@@ -328,22 +539,96 @@ Ils couvrent notamment :
|
||||
- la construction des mesures destinées à la BDD ;
|
||||
- le respect des contraintes du modèle de données.
|
||||
|
||||
Les tests de l'import API Mock couvrent notamment :
|
||||
|
||||
- la récupération des sites ;
|
||||
- l'appel à `/api/v1/readings` ;
|
||||
- les paramètres `site_id`, `start_time`, `end_time` et `limit` ;
|
||||
- la gestion des erreurs HTTP ;
|
||||
- la validation du format de la réponse ;
|
||||
- la transformation des mesures ;
|
||||
- la conservation des valeurs `NULL` ;
|
||||
- la conservation de `data_quality` et `null_reasons` ;
|
||||
- `source = "api_history"` ;
|
||||
- `dataset_id = NULL` ;
|
||||
- la conservation de `raw_data` ;
|
||||
- l'idempotence du chargement.
|
||||
|
||||
Exécuter les tests ETL :
|
||||
|
||||
```powershell
|
||||
uv run pytest tests\etl -v
|
||||
```
|
||||
|
||||
Exécuter les tests unitaires de l'import API Mock :
|
||||
|
||||
```powershell
|
||||
uv run pytest tests\etl\test_mock_api_import.py -v
|
||||
```
|
||||
|
||||
Exécuter le test d'intégration de l'import API Mock :
|
||||
|
||||
```powershell
|
||||
uv run pytest tests\etl\test_mock_api_import.py -m integration -v
|
||||
```
|
||||
|
||||
Contrôler la qualité du code :
|
||||
|
||||
```powershell
|
||||
uv run ruff check app\etl tests\etl
|
||||
```
|
||||
|
||||
## Suite du pipeline Data
|
||||
Contrôler le typage :
|
||||
|
||||
L'import historique constitue la première brique du pipeline Data EnerVision.
|
||||
```powershell
|
||||
uv run mypy app
|
||||
```
|
||||
|
||||
La prochaine étape consiste à orchestrer les traitements ETL avec Apache Airflow, puis à préparer les données nécessaires à l'entraînement du modèle de Machine Learning.
|
||||
Exécuter la suite complète avec le seuil de couverture :
|
||||
|
||||
Airflow sera utilisé comme orchestrateur des traitements existants et ne remplacera pas la logique métier déjà implémentée dans le pipeline ETL.
|
||||
```powershell
|
||||
uv run pytest --cov-fail-under=85
|
||||
```
|
||||
|
||||
Lors de la validation de l'import API Mock :
|
||||
|
||||
```text
|
||||
8 tests unitaires passés
|
||||
1 test d'intégration passé
|
||||
```
|
||||
|
||||
La suite backend complète a également été validée avec une couverture supérieure au seuil de 85 %.
|
||||
|
||||
# Suite du pipeline Data
|
||||
|
||||
Deux sources de données sont maintenant prises en charge :
|
||||
|
||||
```text
|
||||
Dataset CSV/JSON
|
||||
|
|
||||
v
|
||||
historical_import.py
|
||||
|
|
||||
+-----------------+
|
||||
|
|
||||
v
|
||||
PostgreSQL / TimescaleDB
|
||||
^
|
||||
|
|
||||
+-----------------+
|
||||
|
|
||||
mock_api_import.py
|
||||
^
|
||||
|
|
||||
API Mock
|
||||
```
|
||||
|
||||
La logique d'extraction, de transformation et de chargement est donc disponible pour les deux sources de données du MVP.
|
||||
|
||||
La prochaine étape consiste à orchestrer ces traitements avec Apache Airflow.
|
||||
|
||||
Airflow permettra de planifier les traitements, gérer leur ordre d'exécution, suivre leur état et remonter les erreurs.
|
||||
|
||||
Airflow ne remplacera pas la logique ETL Python existante. Les scripts actuels resteront responsables de l'extraction, de la validation, de la transformation et du chargement.
|
||||
|
||||
Le pipeline Data servira ensuite à préparer les données nécessaires au modèle de Machine Learning.
|
||||
Reference in New Issue
Block a user