Compare commits
8
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
56c6b79a5a | ||
|
|
452cfdef85 | ||
|
|
96dd1f834c | ||
|
|
e66ef86729 | ||
|
|
0318ee6cc5 | ||
|
|
0ddfb1997d | ||
|
|
b96546cea3 | ||
|
|
edd5e82d29 |
@@ -17,3 +17,9 @@ APP_LOG_LEVEL=INFO
|
|||||||
APP_SECRET_KEY=change_me
|
APP_SECRET_KEY=change_me
|
||||||
APP_CORS_ORIGINS=http://localhost:4200
|
APP_CORS_ORIGINS=http://localhost:4200
|
||||||
BACKEND_PORT=8000
|
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_PORT=1025
|
||||||
APP_SMTP_USE_TLS=false
|
APP_SMTP_USE_TLS=false
|
||||||
APP_SMTP_FROM_ADDRESS=no-reply@enervision.fr
|
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_pool_size: int = 5
|
||||||
database_max_overflow: int = 10
|
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_issuer: str = "enervision-api"
|
||||||
jwt_audience: str = "enervision-web"
|
jwt_audience: str = "enervision-web"
|
||||||
access_token_ttl_seconds: int = Field(default=900, ge=60, le=3600)
|
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",
|
"argon2-cffi>=23.1",
|
||||||
"anyio>=4.0",
|
"anyio>=4.0",
|
||||||
"aiosmtplib>=5.1.3",
|
"aiosmtplib>=5.1.3",
|
||||||
|
"httpx>=0.28.1",
|
||||||
"pandas>=3.0.5",
|
"pandas>=3.0.5",
|
||||||
]
|
]
|
||||||
|
|
||||||
@@ -27,7 +28,6 @@ dev = [
|
|||||||
"pytest>=9.1.1",
|
"pytest>=9.1.1",
|
||||||
"pytest-asyncio>=1.4.0",
|
"pytest-asyncio>=1.4.0",
|
||||||
"pytest-cov>=7.1.0",
|
"pytest-cov>=7.1.0",
|
||||||
"httpx>=0.28.1",
|
|
||||||
"pandas-stubs>=3.0.5.260914",
|
"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 = "argon2-cffi" },
|
||||||
{ name = "asyncpg" },
|
{ name = "asyncpg" },
|
||||||
{ name = "fastapi" },
|
{ name = "fastapi" },
|
||||||
|
{ name = "httpx" },
|
||||||
{ name = "pandas" },
|
{ name = "pandas" },
|
||||||
{ name = "prometheus-fastapi-instrumentator" },
|
{ name = "prometheus-fastapi-instrumentator" },
|
||||||
{ name = "pydantic", extra = ["email"] },
|
{ name = "pydantic", extra = ["email"] },
|
||||||
@@ -338,7 +339,6 @@ dependencies = [
|
|||||||
|
|
||||||
[package.dev-dependencies]
|
[package.dev-dependencies]
|
||||||
dev = [
|
dev = [
|
||||||
{ name = "httpx" },
|
|
||||||
{ name = "mypy" },
|
{ name = "mypy" },
|
||||||
{ name = "pandas-stubs" },
|
{ name = "pandas-stubs" },
|
||||||
{ name = "pytest" },
|
{ name = "pytest" },
|
||||||
@@ -355,6 +355,7 @@ requires-dist = [
|
|||||||
{ name = "argon2-cffi", specifier = ">=23.1" },
|
{ name = "argon2-cffi", specifier = ">=23.1" },
|
||||||
{ name = "asyncpg", specifier = ">=0.31.0" },
|
{ name = "asyncpg", specifier = ">=0.31.0" },
|
||||||
{ name = "fastapi", specifier = ">=0.141.1" },
|
{ name = "fastapi", specifier = ">=0.141.1" },
|
||||||
|
{ name = "httpx", specifier = ">=0.28.1" },
|
||||||
{ name = "pandas", specifier = ">=3.0.5" },
|
{ name = "pandas", specifier = ">=3.0.5" },
|
||||||
{ name = "prometheus-fastapi-instrumentator", specifier = ">=8.1.0" },
|
{ name = "prometheus-fastapi-instrumentator", specifier = ">=8.1.0" },
|
||||||
{ name = "pydantic", extras = ["email"], specifier = ">=2.13.5" },
|
{ name = "pydantic", extras = ["email"], specifier = ">=2.13.5" },
|
||||||
@@ -367,7 +368,6 @@ requires-dist = [
|
|||||||
|
|
||||||
[package.metadata.requires-dev]
|
[package.metadata.requires-dev]
|
||||||
dev = [
|
dev = [
|
||||||
{ name = "httpx", specifier = ">=0.28.1" },
|
|
||||||
{ name = "mypy", specifier = ">=2.3.1" },
|
{ name = "mypy", specifier = ">=2.3.1" },
|
||||||
{ name = "pandas-stubs", specifier = ">=3.0.5.260914" },
|
{ name = "pandas-stubs", specifier = ">=3.0.5.260914" },
|
||||||
{ name = "pytest", specifier = ">=9.1.1" },
|
{ name = "pytest", specifier = ">=9.1.1" },
|
||||||
|
|||||||
@@ -50,6 +50,12 @@ services:
|
|||||||
APP_SECRET_KEY: ${APP_SECRET_KEY:?}
|
APP_SECRET_KEY: ${APP_SECRET_KEY:?}
|
||||||
APP_CORS_ORIGINS: ${APP_CORS_ORIGINS:-http://localhost:4200}
|
APP_CORS_ORIGINS: ${APP_CORS_ORIGINS:-http://localhost:4200}
|
||||||
DATABASE_URL: postgresql+asyncpg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
|
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_FRONTEND_RESET_PASSWORD_URL: ${APP_FRONTEND_RESET_PASSWORD_URL:-http://localhost:4200/reset-password}
|
||||||
APP_SMTP_HOST: mailpit
|
APP_SMTP_HOST: mailpit
|
||||||
APP_SMTP_PORT: "1025"
|
APP_SMTP_PORT: "1025"
|
||||||
|
|||||||
+332
-69
@@ -6,10 +6,14 @@ système qui en découle.
|
|||||||
|
|
||||||
## Ce que couvre ce document
|
## Ce que couvre ce document
|
||||||
|
|
||||||
**Dix tables applicatives existent** : quatre pour l'authentification, six pour les données
|
**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 le code. Celles
|
d'énergie, dont l'hypertable `reading`.
|
||||||
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.
|
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
|
## Trois emplacements, trois rôles
|
||||||
|
|
||||||
@@ -35,8 +39,16 @@ Statut : `Fait`.
|
|||||||
- `db/init/100-extensions.sql` crée l'extension `timescaledb`.
|
- `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
|
- `db/init/110-test-database.sql` crée `enervision_test`, dont le nom est attendu en dur par
|
||||||
`apps/backend/tests/conftest.py`.
|
`apps/backend/tests/conftest.py`.
|
||||||
- Cinq révisions Alembic. La première, `5353c0e4f094`, **ne crée aucune table** : elle
|
- Six révisions Alembic sont actuellement appliquées.
|
||||||
établit `alembic_version` et refuse de s'appliquer si l'extension manque :
|
- 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
|
```sql
|
||||||
IF NOT EXISTS (SELECT 1 FROM pg_extension WHERE extname = 'timescaledb') THEN
|
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
|
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.
|
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
|
## Cycle de vie d'une mesure
|
||||||
|
|
||||||
Statut : `Cible`, sauf l'hypertable `reading` qui existe. Ni l'ingestion, ni les agrégats
|
Statut : `Partiellement fait`.
|
||||||
continus, ni la compression, ni la rétention ne sont écrits.
|
|
||||||
|
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
|
```mermaid
|
||||||
flowchart LR
|
flowchart LR
|
||||||
src["Source de mesures"] -.-> ing["Ingestion Airflow"]
|
csv["CSV + JSON"] --> hist["historical_import.py"]
|
||||||
ing -.-> hy[("Hypertable reading")]
|
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 -.-> agg[("Agrégat continu")]
|
||||||
hy -.-> comp["Compression"]
|
hy -.-> comp["Compression"]
|
||||||
hy -.-> ret["Rétention"]
|
hy -.-> ret["Rétention"]
|
||||||
agg -.-> api["API FastAPI"]
|
|
||||||
|
agg -.-> backend["API FastAPI"]
|
||||||
agg -.-> graf["Grafana"]
|
agg -.-> graf["Grafana"]
|
||||||
```
|
```
|
||||||
|
|
||||||
Les lectures de l'API et de Grafana visent l'agrégat continu, pas la table brute : c'est tout
|
Les flèches pleines représentent les traitements actuellement implémentés.
|
||||||
l'intérêt de TimescaleDB, et cela doit rester vrai quand les volumes augmenteront.
|
|
||||||
|
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
|
## Tables d'authentification
|
||||||
|
|
||||||
Statut : `Fait`. Elles ne sont pas des séries temporelles et n'ont donc rien à voir avec les
|
Statut : `Fait`.
|
||||||
hypertables ; elles vivent dans `apps/backend/alembic/`, qui porte le schéma exposé par l'API.
|
|
||||||
|
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
|
```mermaid
|
||||||
erDiagram
|
erDiagram
|
||||||
APP_USER ||--o{ REFRESH_TOKEN : ouvre
|
APP_USER ||--o{ REFRESH_TOKEN : ouvre
|
||||||
|
APP_USER ||--o{ PASSWORD_RESET_TOKEN : recoit
|
||||||
|
|
||||||
APP_USER {
|
APP_USER {
|
||||||
uuid id PK
|
uuid id PK
|
||||||
string email UK
|
string email UK
|
||||||
@@ -88,6 +121,7 @@ erDiagram
|
|||||||
bool must_change_password
|
bool must_change_password
|
||||||
timestamptz credentials_changed_at
|
timestamptz credentials_changed_at
|
||||||
}
|
}
|
||||||
|
|
||||||
REFRESH_TOKEN {
|
REFRESH_TOKEN {
|
||||||
uuid id PK
|
uuid id PK
|
||||||
uuid family_id
|
uuid family_id
|
||||||
@@ -99,6 +133,7 @@ erDiagram
|
|||||||
text revoked_reason
|
text revoked_reason
|
||||||
uuid replaced_by
|
uuid replaced_by
|
||||||
}
|
}
|
||||||
|
|
||||||
LOGIN_ATTEMPT {
|
LOGIN_ATTEMPT {
|
||||||
bigint id PK
|
bigint id PK
|
||||||
timestamptz occurred_at
|
timestamptz occurred_at
|
||||||
@@ -106,6 +141,7 @@ erDiagram
|
|||||||
inet client_ip
|
inet client_ip
|
||||||
text outcome
|
text outcome
|
||||||
}
|
}
|
||||||
|
|
||||||
AUDIT_LOG {
|
AUDIT_LOG {
|
||||||
bigint id PK
|
bigint id PK
|
||||||
timestamptz occurred_at
|
timestamptz occurred_at
|
||||||
@@ -114,31 +150,47 @@ erDiagram
|
|||||||
text action
|
text action
|
||||||
jsonb detail
|
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
|
- **`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
|
`CURRENT_USER`. Le nom rappelle aussi qu'il s'agit d'un compte applicatif.
|
||||||
au rôle PostgreSQL qui portera le cantonnement de l'ETL.
|
- **`credentials_changed_at`, une seule colonne**, couvre notamment le changement de mot de passe,
|
||||||
- **`credentials_changed_at`, une seule colonne**, couvre le changement de mot de passe, le
|
le changement de rôle et la désactivation.
|
||||||
changement de rôle et la désactivation. Un compteur de version ne dirait rien à un humain qui
|
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation.
|
||||||
lit un audit.
|
- **`audit_log.actor_id` n'a aucune clé étrangère** afin de conserver les informations d'audit
|
||||||
- **`refresh_token.expires_at` est absolu et hérité** du prédécesseur à chaque rotation. S'il
|
même si l'entité d'origine évolue.
|
||||||
glissait, la promesse de sept jours serait fictive et une session active ne finirait jamais.
|
- `password_reset_token` ne stocke que l'empreinte du jeton et jamais sa valeur directement.
|
||||||
- **`audit_log.actor_id` n'a aucune clé étrangère**, et `actor_email` comme `actor_role` sont
|
- `password_reset_attempt` est séparée de `audit_log`, car son volume peut être piloté
|
||||||
dénormalisés. Une contrainte `ON DELETE SET NULL` déclencherait un `UPDATE` que le déclencheur
|
par des demandes externes répétées.
|
||||||
d'ajout seul refuserait. Voir l'[ADR 0004](../adr/0004-journal-d-audit-en-ajout-seul.md).
|
|
||||||
|
|
||||||
`audit_log` porte deux déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`. Elle n'est
|
`audit_log` porte des déclencheurs qui refusent `UPDATE`, `DELETE` et `TRUNCATE`.
|
||||||
donc **pas** une hypertable : une politique de rétention émettrait des `DELETE` qu'ils
|
Elle n'est donc **pas** une hypertable.
|
||||||
refuseraient. `login_attempt`, à l'inverse, est faite pour se purger, puisque son volume est
|
|
||||||
piloté par l'attaquant.
|
|
||||||
|
|
||||||
## Gabarit de révision créant 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
|
Conforme à la règle de l'ADR 0001 : table et hypertable dans la même révision.
|
||||||
`e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
|
||||||
|
La révision `e6d2026091501` en est l'exemple réel, réduit ici à l'essentiel.
|
||||||
|
|
||||||
```python
|
```python
|
||||||
def upgrade() -> None:
|
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
|
## Questions ouvertes
|
||||||
|
|
||||||
Elles relèvent du jalon J2, « valider le périmètre retenu ». Le schéma est livré : ce qui suit
|
Elles portent maintenant principalement sur l'exploitation du schéma :
|
||||||
porte sur son exploitation, plus sur sa forme.
|
|
||||||
|
|
||||||
- **Quelle granularité** à l'ingestion : la seconde, la minute, le quart d'heure.
|
- **Quelle granularité** conserver à long terme à l'ingestion : seconde, minute ou quart d'heure.
|
||||||
- **Quels agrégats continus**, et sur quelles fenêtres.
|
- **Quels agrégats continus** créer et sur quelles fenêtres.
|
||||||
- **Quelle profondeur de rétention** en données brutes, et à partir de quand on compresse.
|
- **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.
|
- **Multi-tenant ou non** : un site appartient-il à un client et faut-il cloisonner les lectures.
|
||||||
|
|
||||||
## Modélisation détaillée des données
|
## Modélisation détaillée des données
|
||||||
|
|
||||||
Cette modélisation prend en compte les fichiers CSV historiques,
|
Cette modélisation prend en compte :
|
||||||
leurs métadonnées JSON et les données de l’API Mock.
|
|
||||||
Elle comprend six tables, depuis le stockage des mesures
|
- les fichiers CSV historiques ;
|
||||||
jusqu’aux recommandations proposées à l’utilisateur.
|
- 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
|
### Schéma de données
|
||||||
|
|
||||||
Le diagramme ci-dessous présente les tables et leurs relations.
|
Le diagramme ci-dessous présente les tables et leurs relations.
|
||||||
|
|
||||||
La révision `e6d2026091501` les crée.
|
La révision `e6d2026091501` les crée.
|
||||||
|
|
||||||

|

|
||||||
@@ -208,21 +264,20 @@ La révision `e6d2026091501` les crée.
|
|||||||
|
|
||||||
### Description des tables
|
### Description des tables
|
||||||
|
|
||||||
Chaque table remplit un rôle précis dans le traitement et l’exploitation
|
Chaque table remplit un rôle précis dans le traitement et l'exploitation des données.
|
||||||
des données.
|
|
||||||
|
|
||||||
| Table | Rôle | Origine des informations |
|
| 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` |
|
| `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` |
|
| `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 |
|
| `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
|
Les anomalies historiques décrites dans les JSON sont conservées dans `dataset.metadata`.
|
||||||
dans `dataset.metadata`. Elles servent à l’analyse des données
|
|
||||||
et ne sont pas considérées comme des alertes actuelles.
|
Elles servent à l'analyse des données et ne sont pas considérées comme des alertes actuelles.
|
||||||
|
|
||||||
### Relations entre les tables
|
### 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 être associée à une prévision du même site.
|
||||||
- Une alerte peut donner lieu à plusieurs recommandations.
|
- 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
|
```text
|
||||||
Dataset CSV + métadonnées JSON
|
Dataset CSV + métadonnées JSON
|
||||||
@@ -271,17 +332,26 @@ Dataset CSV + métadonnées JSON
|
|||||||
|
|
||||||
Le pipeline est développé en Python.
|
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é.
|
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.
|
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 :
|
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 ;
|
- 122 647 mesures ;
|
||||||
- 0 doublon détecté dans le dataset source.
|
- 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
|
## 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 :
|
Le pipeline assure :
|
||||||
|
|
||||||
@@ -12,12 +13,15 @@ Le pipeline assure :
|
|||||||
- la validation de leur structure et de leur cohérence ;
|
- la validation de leur structure et de leur cohérence ;
|
||||||
- la normalisation des données nécessaires au stockage ;
|
- la normalisation des données nécessaires au stockage ;
|
||||||
- le suivi de la qualité des données ;
|
- 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 ;
|
- 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.
|
- l'idempotence du chargement afin d'éviter la création de doublons.
|
||||||
|
|
||||||
## Données sources
|
## Données sources
|
||||||
|
|
||||||
|
### Dataset historique
|
||||||
|
|
||||||
Le dataset est fourni par le formateur dans le cadre du projet EnerVision.
|
Le dataset est fourni par le formateur dans le cadre du projet EnerVision.
|
||||||
|
|
||||||
Il contient les deux fichiers suivants :
|
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.
|
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
|
```text
|
||||||
data/raw/
|
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.
|
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
|
## Technologies utilisées
|
||||||
|
|
||||||
| Technologie | Utilisation |
|
| Technologie | Utilisation |
|
||||||
|---|---|
|
|---|---|
|
||||||
| Python | Développement du pipeline ETL |
|
| Python | Développement du pipeline ETL |
|
||||||
| Pandas | Lecture, validation et transformation des données |
|
| Pandas | Lecture, validation et transformation du dataset historique |
|
||||||
| JSON | Lecture des métadonnées du dataset |
|
| JSON | Lecture des métadonnées et conservation des données sources |
|
||||||
| hashlib / SHA-256 | Identification, intégrité et traçabilité du dataset |
|
| 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 |
|
| SQLAlchemy Async | Connexion et chargement asynchrone en base |
|
||||||
| PostgreSQL | Stockage relationnel |
|
| PostgreSQL | Stockage relationnel |
|
||||||
| TimescaleDB | Stockage des séries temporelles énergétiques |
|
| 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 |
|
| Alembic | Gestion des migrations du schéma |
|
||||||
| uv | Gestion et exécution de l'environnement Python |
|
| uv | Gestion et exécution de l'environnement Python |
|
||||||
| Ruff | Contrôle de la qualité du code |
|
| Ruff | Contrôle de la qualité du code |
|
||||||
|
| mypy | Vérification du typage |
|
||||||
| Pytest | Tests automatisés |
|
| 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
|
```text
|
||||||
apps/backend/app/etl/historical_import.py
|
apps/backend/app/etl/historical_import.py
|
||||||
@@ -207,7 +226,7 @@ Valeurs manquantes identifiées :
|
|||||||
| `humidity_percent` | 3 423 |
|
| `humidity_percent` | 3 423 |
|
||||||
| `solar_irradiance_wm2` | 3 964 |
|
| `solar_irradiance_wm2` | 3 964 |
|
||||||
|
|
||||||
## Exécution en dry-run
|
## Exécution historique en dry-run
|
||||||
|
|
||||||
Depuis le dossier :
|
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.
|
Aucune donnée n'est écrite dans la base pendant cette exécution.
|
||||||
|
|
||||||
## Chargement réel
|
## Chargement historique réel
|
||||||
|
|
||||||
Depuis `apps/backend/` :
|
Depuis `apps/backend/` :
|
||||||
|
|
||||||
@@ -249,7 +268,7 @@ Chargement : 2000/122647
|
|||||||
Chargement : 122647/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é :
|
Après le chargement initial, les contrôles en base ont confirmé :
|
||||||
|
|
||||||
@@ -266,7 +285,7 @@ Le premier import a créé :
|
|||||||
nouvelles lectures : 122647
|
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.
|
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é.
|
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 :
|
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;"
|
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
|
```text
|
||||||
csv | 122647
|
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
|
```text
|
||||||
apps/backend/tests/etl/
|
apps/backend/tests/etl/
|
||||||
```
|
```
|
||||||
|
|
||||||
Ils couvrent notamment :
|
Les tests de l'import historique couvrent notamment :
|
||||||
|
|
||||||
- la validation du dataset ;
|
- la validation du dataset ;
|
||||||
- les colonnes obligatoires ;
|
- les colonnes obligatoires ;
|
||||||
@@ -328,22 +539,96 @@ Ils couvrent notamment :
|
|||||||
- la construction des mesures destinées à la BDD ;
|
- la construction des mesures destinées à la BDD ;
|
||||||
- le respect des contraintes du modèle de données.
|
- 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 :
|
Exécuter les tests ETL :
|
||||||
|
|
||||||
```powershell
|
```powershell
|
||||||
uv run pytest tests\etl -v
|
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 :
|
Contrôler la qualité du code :
|
||||||
|
|
||||||
```powershell
|
```powershell
|
||||||
uv run ruff check app\etl tests\etl
|
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.
|
||||||
@@ -9,7 +9,7 @@ sonar.tests=apps/frontend/src,apps/backend/tests
|
|||||||
sonar.test.inclusions=**/*.spec.ts,**/*.test.ts,**/*test_*.py,**/*test.py
|
sonar.test.inclusions=**/*.spec.ts,**/*.test.ts,**/*test_*.py,**/*test.py
|
||||||
|
|
||||||
# Liste des fichiers et dossiers à exclure de l'analyse
|
# Liste des fichiers et dossiers à exclure de l'analyse
|
||||||
sonar.exclusions=.pytest_cache,.venv,alembic,tests,**/*/node_modules/**,**/*/dist/**,**/*/build/**,**/*.spec.ts,**/*.test.ts,**/*test_*.py,**/*test.py
|
sonar.exclusions=.pytest_cache,.venv,alembic,tests,**/*/node_modules/**,**/*/dist/**,**/*/build/**,**/*.spec.ts,**/*.test.ts,**/*test_*.py,**/*test.py,**/*.spec.ts
|
||||||
|
|
||||||
# Chemin vers le rapport de couverture de code
|
# Chemin vers le rapport de couverture de code
|
||||||
# Fichier généré par Pytest
|
# Fichier généré par Pytest
|
||||||
|
|||||||
Reference in New Issue
Block a user