Aller au contenu

Pipeline (kadi.kidas.pipeline)

DataPipeline est le chef d'orchestre du module kidas. Il enchaîne les étapes de chargement, nettoyage, validation et normalisation dans un flux de traitement configurable et reproductible.


Principe

Le pipeline suit un pattern de construction en chaîne (fluent interface) : chaque méthode retourne l'objet pipeline lui-même, ce qui permet d'enchaîner les appels de manière lisible.

load_data()
    → add_cleaning_step()
    → add_validation_step()
    → add_normalization_step()
    → execute()
    → (DataFrame, rapport)

Initialisation

from kadi.kidas import DataPipeline

pipeline = DataPipeline()

Méthodes

load_data(source)

Charge les données depuis une source. Le format est détecté automatiquement.

pipeline.load_data("recoltes_2024.csv")         # CSV
pipeline.load_data("marches_prix.xlsx")         # Excel
pipeline.load_data("capteurs_meteo.json")        # JSON
pipeline.load_data("donnees_sol.nc")             # NetCDF
pipeline.load_data("https://api.data.bj/...")   # API REST

Paramètre : source — chemin de fichier ou URL.


add_cleaning_step(step_name, **kwargs)

Ajoute une étape de nettoyage à la file d'exécution.

pipeline.add_cleaning_step("remove_duplicates")
pipeline.add_cleaning_step("handle_missing_values", strategy="median")
pipeline.add_cleaning_step("remove_outliers", method="zscore", threshold=3.0)
pipeline.add_cleaning_step("normalize_text", columns=["culture", "marche"])

Étapes disponibles :

Nom Paramètres Description
remove_duplicates Supprime les lignes dupliquées
handle_missing_values strategy: 'mean', 'median', 'drop' Impute ou supprime les valeurs manquantes
remove_outliers method: 'zscore'/'iqr', threshold Supprime les valeurs aberrantes
normalize_text columns: liste Nettoie accents, majuscules, espaces
fix_encoding Corrige les problèmes d'encodage

add_validation_step(schema)

Ajoute une étape de validation qui vérifie les types et contraintes des colonnes.

pipeline.add_validation_step({
    "culture": "str",         # Doit être une chaîne
    "rendement_kg": "float",  # Doit être un nombre
    "date_recolte": "date",   # Doit être une date
    "latitude": "float",      # Doit être un flottant
})

Types supportés : 'str', 'int', 'float', 'bool', 'date'.

La validation ne bloque pas l'exécution. Les erreurs sont collectées dans le rapport de pipeline sous la clé warnings.


add_normalization_step(mapping)

Ajoute une étape de normalisation qui standardise les noms de colonnes.

pipeline.add_normalization_step({
    "crops": "culture",        # Normalise la colonne "culture" vers les codes KadiPy
    "markets": "ville",        # Normalise les noms de marchés
    "gps": ["latitude", "longitude"],  # Valide et normalise les coordonnées GPS
})

execute(cache)

Exécute toutes les étapes enregistrées et retourne le résultat.

df, rapport = pipeline.execute(cache=True)

Paramètre :

Nom Type Défaut Description
cache bool True Si True, met en cache le résultat dans SQLite

Retour : tuple[pd.DataFrame, dict]


Rapport de pipeline

df, rapport = pipeline.execute()

# Nombre de lignes en entrée et en sortie
print(f"Lignes en entrée : {rapport['nb_rows_in']}")
print(f"Lignes en sortie : {rapport['nb_rows_out']}")
print(f"Lignes supprimées: {rapport['nb_rows_in'] - rapport['nb_rows_out']}")

# Score de qualité global (0 à 1)
qualite = rapport["quality_score"]
print(f"Score global     : {qualite['overall']:.2f}")
print(f"Score complétude : {qualite.get('completeness', 'N/A')}")
print(f"Score cohérence  : {qualite.get('consistency', 'N/A')}")

# Résumé de chaque étape
for etape in rapport["steps_summary"]:
    print(f"  [{etape['step']}] {etape['status']}{etape.get('detail', '')}")

# Avertissements de validation
for avertissement in rapport["warnings"]:
    print(f"  Attention : {avertissement}")

Exemple complet avec toutes les étapes

from kadi.kidas import DataPipeline

pipeline = DataPipeline()

df, rapport = (
    pipeline
    .load_data("enquete_agriculteurs_2024.csv")
    .add_cleaning_step("fix_encoding")
    .add_cleaning_step("remove_duplicates")
    .add_cleaning_step("normalize_text", columns=["culture", "commune"])
    .add_cleaning_step("handle_missing_values", strategy="median")
    .add_cleaning_step("remove_outliers", method="iqr")
    .add_validation_step({
        "culture": "str",
        "rendement_kg_ha": "float",
        "superficie_ha": "float",
        "commune": "str",
        "latitude": "float",
        "longitude": "float",
    })
    .add_normalization_step({
        "crops": "culture",
        "gps": ["latitude", "longitude"],
    })
    .execute(cache=True)
)

print(f"Données prêtes : {len(df)} lignes, {len(df.columns)} colonnes")
print(f"Score qualité  : {rapport['quality_score']['overall']:.2f} / 1.0")

DataPipeline

Orchestrateur du flux de traitement de données agricoles kidas.

DataPipeline est le point d'entrée unique pour les utilisateurs de kidas. Il permet de définir de manière déclarative et chainable un flux complet : chargement → nettoyage → validation → normalisation → cache.

L'auto-détection du type de source (CSV, Excel, JSON, NetCDF, API) est basée sur l'extension du fichier ou le préfixe 'http' de l'URL.

Attributs

_source (DataSource | None): Source de données configurée. _df (pd.DataFrame | None): Données courantes dans le pipeline. _etapes (list): Liste ordonnée des étapes de traitement configurées. _rapports (dict): Rapports agrégés de toutes les étapes. _cache (DataCache): Instance du gestionnaire de cache kidas.

Exemple

pipeline = DataPipeline() df, report = ( ... pipeline ... .load_data('recolte_2024.xlsx') ... .add_cleaning_step('remove_duplicates') ... .add_cleaning_step('handle_missing_values', strategy='mean') ... .add_validation_step({'culture': 'str', 'rendement_kg': 'float'}) ... .add_normalization_step({'culture': 'fao_standard'}) ... .execute(cache=True) ... ) print(report['lignes_finales']) 150

Source code in kadi/kidas/pipeline.py
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
class DataPipeline:
    """Orchestrateur du flux de traitement de données agricoles kidas.

    DataPipeline est le point d'entrée unique pour les utilisateurs de kidas.
    Il permet de définir de manière déclarative et chainable un flux complet :
    chargement → nettoyage → validation → normalisation → cache.

    L'auto-détection du type de source (CSV, Excel, JSON, NetCDF, API) est
    basée sur l'extension du fichier ou le préfixe 'http' de l'URL.

    Attributs:
        _source (DataSource | None): Source de données configurée.
        _df (pd.DataFrame | None): Données courantes dans le pipeline.
        _etapes (list): Liste ordonnée des étapes de traitement configurées.
        _rapports (dict): Rapports agrégés de toutes les étapes.
        _cache (DataCache): Instance du gestionnaire de cache kidas.

    Exemple:
        >>> pipeline = DataPipeline()
        >>> df, report = (
        ...     pipeline
        ...     .load_data('recolte_2024.xlsx')
        ...     .add_cleaning_step('remove_duplicates')
        ...     .add_cleaning_step('handle_missing_values', strategy='mean')
        ...     .add_validation_step({'culture': 'str', 'rendement_kg': 'float'})
        ...     .add_normalization_step({'culture': 'fao_standard'})
        ...     .execute(cache=True)
        ... )
        >>> print(report['lignes_finales'])
        150
    """

    def __init__(self) -> None:
        """Initialise un pipeline vide prêt à recevoir des étapes de traitement."""
        # Source de données (sera configurée par load_data())
        self._source: Optional[DataSource] = None

        # DataFrame courant (None jusqu'à l'appel de execute())
        self._df: Optional[pd.DataFrame] = None

        # Liste ordonnée des étapes : {'type', 'nom', 'params'}
        self._etapes: List[Dict[str, Any]] = []

        # Rapports agrégés des étapes de traitement
        self._rapports: Dict[str, Any] = {
            "source": None,
            "etapes_appliquees": [],
            "nettoyage": None,
            "validation": None,
            "normalisation": None,
            "cache_utilise": False,
        }

        # Instance du cache kidas
        self._cache: DataCache = DataCache()

    @staticmethod
    def _detecter_type_source(source: str) -> str:
        """Détecte le type de source de données à partir du chemin ou de l'URL.

        Args:
            source (str): Chemin vers le fichier ou URL de l'API.

        Returns:
            str: Type détecté parmi 'csv', 'excel', 'json', 'netcdf', 'api'.

        Raises:
            KidasPipelineError: Si le type de source ne peut pas être déterminé.
        """
        # Détection des APIs par préfixe HTTP/HTTPS
        if source.startswith("http://") or source.startswith("https://"):
            return "api"

        # Détection par extension de fichier
        _, extension = os.path.splitext(source.lower())

        if extension in _EXT_CSV:
            return "csv"
        elif extension in _EXT_EXCEL:
            return "excel"
        elif extension in _EXT_JSON:
            return "json"
        elif extension in _EXT_NETCDF:
            return "netcdf"
        else:
            raise KidasPipelineError(
                f"Impossible de détecter le type de source pour '{source}'. "
                f"Extensions supportées : CSV {_EXT_CSV}, Excel {_EXT_EXCEL}, "
                f"JSON {_EXT_JSON}, NetCDF {_EXT_NETCDF}, ou URL HTTP."
            )

    def load_data(
        self,
        source: Union[str, DataSource],
        **kwargs: Any,
    ) -> "DataPipeline":
        """Configure la source de données du pipeline.

        Auto-détecte le type de source si un chemin de fichier est fourni,
        ou utilise directement une instance DataSource existante.

        Args:
            source (str | DataSource): Chemin vers le fichier, URL de l'API,
                ou instance DataSource directement.
            **kwargs: Arguments optionnels transmis au constructeur de la
                DataSource (ex: encoding, sheet_name, use_dask).

        Returns:
            DataPipeline: L'instance courante (pour le chaînage).

        Raises:
            KidasPipelineError: Si le type de source est indéterminable.
        """
        if isinstance(source, DataSource):
            # Utilisation directe d'une DataSource existante
            self._source = source
            type_source = source.source_type
        else:
            # Auto-détection et instanciation de la DataSource appropriée
            type_source = self._detecter_type_source(source)

            if type_source == "csv":
                self._source = CSVDataSource(source, **kwargs)
            elif type_source == "excel":
                self._source = ExcelDataSource(source, **kwargs)
            elif type_source == "json":
                self._source = JSONDataSource(source, **kwargs)
            elif type_source == "netcdf":
                self._source = NetCDFDataSource(source, **kwargs)
            elif type_source == "api":
                self._source = APIDataSource(source, **kwargs)

        # Enregistrement dans le rapport
        self._rapports["source"] = {
            "path": str(source),
            "type": type_source,
        }

        logger.info(
            "Pipeline kidas : source '%s' configurée (type: %s).",
            source,
            type_source,
        )

        return self

    def add_cleaning_step(
        self,
        step_name: str,
        **params: Any,
    ) -> "DataPipeline":
        """Ajoute une étape de nettoyage à la file du pipeline.

        Les étapes sont exécutées dans l'ordre de leur ajout lors de
        l'appel à execute().

        Args:
            step_name (str): Nom de la méthode DataCleaner à appeler.
                Valeurs acceptées : 'remove_duplicates', 'handle_missing_values',
                'remove_outliers', 'fix_dates', 'standardize_text',
                'remove_special_chars'.
            **params: Paramètres à passer à la méthode de nettoyage
                (ex: strategy='mean', method='iqr').

        Returns:
            DataPipeline: L'instance courante (pour le chaînage).
        """
        # Enregistrement de l'étape dans la file
        self._etapes.append({
            "type": "cleaning",
            "nom": step_name,
            "params": params,
        })

        logger.debug(
            "Étape de nettoyage ajoutée : '%s' (params: %s).", step_name, params
        )

        return self

    def add_validation_step(
        self,
        schema: Dict[str, str],
    ) -> "DataPipeline":
        """Ajoute une étape de validation de schéma à la file du pipeline.

        Args:
            schema (dict[str, str]): Schéma de validation : nom_colonne → type.
                Exemple : {'culture': 'str', 'rendement_kg': 'float'}.

        Returns:
            DataPipeline: L'instance courante (pour le chaînage).
        """
        self._etapes.append({
            "type": "validation",
            "nom": "validate_schema",
            "params": {"schema": schema},
        })

        logger.debug("Étape de validation ajoutée (schéma: %s).", schema)

        return self

    def add_normalization_step(
        self,
        mappings: Dict[str, Any],
    ) -> "DataPipeline":
        """Ajoute une étape de normalisation à la file du pipeline.

        Args:
            mappings (dict): Dictionnaire de normalisation. Les clés supportées
                sont 'columns' (snake_case), 'units' (unités), 'crops' (FAO),
                'markets' (Bénin). Exemple : {'crops': 'culture'}.

        Returns:
            DataPipeline: L'instance courante (pour le chaînage).
        """
        self._etapes.append({
            "type": "normalization",
            "nom": "normalize",
            "params": {"mappings": mappings},
        })

        logger.debug("Étape de normalisation ajoutée (mappings: %s).", mappings)

        return self

    def execute(
        self,
        cache: bool = True,
    ) -> Tuple[pd.DataFrame, Dict]:
        """Exécute toutes les étapes configurées du pipeline.

        Charge les données depuis la source, applique les étapes de
        nettoyage, validation et normalisation dans l'ordre, puis
        met le résultat en cache si demandé.

        Args:
            cache (bool): Si True, tente de charger depuis le cache avant
                la lecture et sauvegarde le résultat final. Par défaut True.

        Returns:
            tuple[pd.DataFrame, dict]: Tuple contenant :
                - Le DataFrame traité et prêt à l'emploi.
                - Le rapport complet de toutes les étapes.

        Raises:
            KidasPipelineError: Si aucune source n'a été configurée.
            KidasReadError: Si la lecture de la source échoue.
        """
        # Vérification qu'une source a été configurée
        if self._source is None:
            raise KidasPipelineError(
                "Aucune source configurée. Appelez load_data() avant execute()."
            )

        # Génération d'une clé de cache basée sur le chemin de la source
        cle_cache = f"pipeline_{self._source.source_path}"

        # Tentative de chargement depuis le cache
        if cache:
            df_cached, _ = self._cache.load(cle_cache)
            if df_cached is not None:
                logger.info(
                    "Pipeline kidas : données chargées depuis le cache (clé: '%s').",
                    cle_cache,
                )
                self._rapports["cache_utilise"] = True
                return df_cached, self._rapports

        # Lecture des données depuis la source
        try:
            self._df = self._source.read()
            logger.info(
                "Pipeline kidas : %d lignes chargées depuis '%s'.",
                len(self._df),
                self._source.source_path,
            )
        except Exception as erreur:
            raise KidasReadError(
                f"Échec de lecture dans le pipeline : {erreur}"
            ) from erreur

        # Exécution des étapes dans l'ordre
        for etape in self._etapes:
            self._df = self._executer_etape(etape)

        # Sauvegarde en cache du résultat final
        if cache and self._df is not None:
            self._cache.save(cle_cache, self._df)
            self._rapports["cache_utilise"] = True

        # Compilation du rapport final
        self._rapports["lignes_finales"] = len(self._df) if self._df is not None else 0
        self._rapports["etapes_appliquees"] = [e["nom"] for e in self._etapes]

        return self._df, self._rapports

    def _executer_etape(self, etape: Dict) -> pd.DataFrame:
        """Exécute une étape individuelle du pipeline sur le DataFrame courant.

        Args:
            etape (dict): Dictionnaire décrivant l'étape :
                {'type', 'nom', 'params'}.

        Returns:
            pd.DataFrame: Le DataFrame résultant de l'étape.

        Raises:
            KidasPipelineError: Si la méthode de l'étape est inconnue.
        """
        type_etape = etape["type"]
        nom_methode = etape["nom"]
        params = etape["params"]

        try:
            if type_etape == "cleaning":
                # Instanciation du nettoyeur et appel dynamique de la méthode
                cleaner = DataCleaner(self._df)

                if not hasattr(cleaner, nom_methode):
                    raise KidasPipelineError(
                        f"Méthode de nettoyage '{nom_methode}' inconnue. "
                        f"Méthodes disponibles : remove_duplicates, "
                        f"handle_missing_values, remove_outliers, fix_dates, "
                        f"standardize_text, remove_special_chars."
                    )

                # Appel de la méthode avec les paramètres
                resultat = getattr(cleaner, nom_methode)(**params)

                # Certaines méthodes retournent un tuple (df, outliers_df)
                if isinstance(resultat, tuple):
                    self._df = resultat[0]
                else:
                    self._df = resultat

                # Mise à jour du rapport de nettoyage
                self._rapports["nettoyage"] = cleaner.get_cleaning_report()

            elif type_etape == "validation":
                # Validation du schéma et calcul du score qualité
                validator = DataValidator(self._df)
                est_valide, erreurs = validator.validate_schema(params["schema"])
                score = validator.compute_quality_score()

                # Journalisation du résultat de validation
                if not est_valide:
                    logger.warning(
                        "Validation du schéma : %d erreur(s) détectée(s).", len(erreurs)
                    )
                    for erreur in erreurs:
                        logger.warning("  - %s", erreur)

                self._rapports["validation"] = validator.get_validation_report()
                self._rapports["quality_score"] = score

            elif type_etape == "normalization":
                # Application des normalisations demandées
                normalizer = DataNormalizer(self._df)
                mappings = params["mappings"]

                # Normalisation des noms de colonnes si demandé
                if "columns" in mappings or mappings.get("normalize_columns"):
                    normalizer.normalize_column_names()

                # Normalisation des noms de cultures si demandé
                if "crops" in mappings:
                    normalizer.normalize_crop_names(col=mappings["crops"])

                # Normalisation des unités si demandé
                if "units" in mappings:
                    normalizer.normalize_units(unit_map=mappings["units"])

                # Normalisation des marchés si demandé
                if "markets" in mappings:
                    normalizer.normalize_market_names(col=mappings["markets"])

                self._df = normalizer.df
                self._rapports["normalisation"] = normalizer.get_normalization_mapping()

        except KidasPipelineError:
            raise
        except Exception as erreur:
            raise KidasPipelineError(
                f"Erreur lors de l'exécution de l'étape '{nom_methode}' : {erreur}"
            ) from erreur

        return self._df

    def get_pipeline_config(self) -> dict:
        """Retourne la configuration complète du pipeline (étapes définies).

        Returns:
            dict: Dictionnaire décrivant la source et les étapes configurées.
        """
        return {
            "source": self._rapports.get("source"),
            "nb_etapes": len(self._etapes),
            "etapes": [
                {"type": e["type"], "nom": e["nom"], "params": e["params"]}
                for e in self._etapes
            ],
        }

    def export_report(self, filepath: str) -> bool:
        """Exporte le rapport de pipeline dans un fichier JSON ou HTML.

        Args:
            filepath (str): Chemin de destination du rapport. L'extension
                détermine le format : '.json' ou '.html'.

        Returns:
            bool: True si l'export s'est déroulé avec succès.

        Raises:
            KidasPipelineError: Si l'extension n'est pas supportée.
        """
        import json

        _, extension = os.path.splitext(filepath.lower())

        try:
            if extension == ".json":
                # Export au format JSON
                with open(filepath, "w", encoding="utf-8") as f:
                    json.dump(self._rapports, f, ensure_ascii=False, indent=2, default=str)

            elif extension == ".html":
                # Export au format HTML simplifié
                contenu_json = json.dumps(
                    self._rapports, ensure_ascii=False, indent=2, default=str
                )
                html = (
                    "<html><head><meta charset='utf-8'>"
                    "<title>Rapport kidas Pipeline</title></head>"
                    "<body><h1>Rapport kidas DataPipeline</h1>"
                    f"<pre>{contenu_json}</pre></body></html>"
                )
                with open(filepath, "w", encoding="utf-8") as f:
                    f.write(html)

            else:
                raise KidasPipelineError(
                    f"Format d'export '{extension}' non supporté. "
                    f"Utilisez '.json' ou '.html'."
                )

            logger.info("Rapport pipeline exporté vers '%s'.", filepath)
            return True

        except OSError as erreur:
            raise KidasPipelineError(
                f"Impossible d'écrire le rapport vers '{filepath}' : {erreur}"
            ) from erreur

__init__

__init__() -> None

Initialise un pipeline vide prêt à recevoir des étapes de traitement.

Source code in kadi/kidas/pipeline.py
def __init__(self) -> None:
    """Initialise un pipeline vide prêt à recevoir des étapes de traitement."""
    # Source de données (sera configurée par load_data())
    self._source: Optional[DataSource] = None

    # DataFrame courant (None jusqu'à l'appel de execute())
    self._df: Optional[pd.DataFrame] = None

    # Liste ordonnée des étapes : {'type', 'nom', 'params'}
    self._etapes: List[Dict[str, Any]] = []

    # Rapports agrégés des étapes de traitement
    self._rapports: Dict[str, Any] = {
        "source": None,
        "etapes_appliquees": [],
        "nettoyage": None,
        "validation": None,
        "normalisation": None,
        "cache_utilise": False,
    }

    # Instance du cache kidas
    self._cache: DataCache = DataCache()

load_data

load_data(source: Union[str, DataSource], **kwargs: Any) -> DataPipeline

Configure la source de données du pipeline.

Auto-détecte le type de source si un chemin de fichier est fourni, ou utilise directement une instance DataSource existante.

Parameters:

Name Type Description Default
source str | DataSource

Chemin vers le fichier, URL de l'API, ou instance DataSource directement.

required
**kwargs Any

Arguments optionnels transmis au constructeur de la DataSource (ex: encoding, sheet_name, use_dask).

{}

Returns:

Name Type Description
DataPipeline DataPipeline

L'instance courante (pour le chaînage).

Raises:

Type Description
KidasPipelineError

Si le type de source est indéterminable.

Source code in kadi/kidas/pipeline.py
def load_data(
    self,
    source: Union[str, DataSource],
    **kwargs: Any,
) -> "DataPipeline":
    """Configure la source de données du pipeline.

    Auto-détecte le type de source si un chemin de fichier est fourni,
    ou utilise directement une instance DataSource existante.

    Args:
        source (str | DataSource): Chemin vers le fichier, URL de l'API,
            ou instance DataSource directement.
        **kwargs: Arguments optionnels transmis au constructeur de la
            DataSource (ex: encoding, sheet_name, use_dask).

    Returns:
        DataPipeline: L'instance courante (pour le chaînage).

    Raises:
        KidasPipelineError: Si le type de source est indéterminable.
    """
    if isinstance(source, DataSource):
        # Utilisation directe d'une DataSource existante
        self._source = source
        type_source = source.source_type
    else:
        # Auto-détection et instanciation de la DataSource appropriée
        type_source = self._detecter_type_source(source)

        if type_source == "csv":
            self._source = CSVDataSource(source, **kwargs)
        elif type_source == "excel":
            self._source = ExcelDataSource(source, **kwargs)
        elif type_source == "json":
            self._source = JSONDataSource(source, **kwargs)
        elif type_source == "netcdf":
            self._source = NetCDFDataSource(source, **kwargs)
        elif type_source == "api":
            self._source = APIDataSource(source, **kwargs)

    # Enregistrement dans le rapport
    self._rapports["source"] = {
        "path": str(source),
        "type": type_source,
    }

    logger.info(
        "Pipeline kidas : source '%s' configurée (type: %s).",
        source,
        type_source,
    )

    return self

add_cleaning_step

add_cleaning_step(step_name: str, **params: Any) -> DataPipeline

Ajoute une étape de nettoyage à la file du pipeline.

Les étapes sont exécutées dans l'ordre de leur ajout lors de l'appel à execute().

Parameters:

Name Type Description Default
step_name str

Nom de la méthode DataCleaner à appeler. Valeurs acceptées : 'remove_duplicates', 'handle_missing_values', 'remove_outliers', 'fix_dates', 'standardize_text', 'remove_special_chars'.

required
**params Any

Paramètres à passer à la méthode de nettoyage (ex: strategy='mean', method='iqr').

{}

Returns:

Name Type Description
DataPipeline DataPipeline

L'instance courante (pour le chaînage).

Source code in kadi/kidas/pipeline.py
def add_cleaning_step(
    self,
    step_name: str,
    **params: Any,
) -> "DataPipeline":
    """Ajoute une étape de nettoyage à la file du pipeline.

    Les étapes sont exécutées dans l'ordre de leur ajout lors de
    l'appel à execute().

    Args:
        step_name (str): Nom de la méthode DataCleaner à appeler.
            Valeurs acceptées : 'remove_duplicates', 'handle_missing_values',
            'remove_outliers', 'fix_dates', 'standardize_text',
            'remove_special_chars'.
        **params: Paramètres à passer à la méthode de nettoyage
            (ex: strategy='mean', method='iqr').

    Returns:
        DataPipeline: L'instance courante (pour le chaînage).
    """
    # Enregistrement de l'étape dans la file
    self._etapes.append({
        "type": "cleaning",
        "nom": step_name,
        "params": params,
    })

    logger.debug(
        "Étape de nettoyage ajoutée : '%s' (params: %s).", step_name, params
    )

    return self

add_validation_step

add_validation_step(schema: Dict[str, str]) -> DataPipeline

Ajoute une étape de validation de schéma à la file du pipeline.

Parameters:

Name Type Description Default
schema dict[str, str]

Schéma de validation : nom_colonne → type. Exemple : {'culture': 'str', 'rendement_kg': 'float'}.

required

Returns:

Name Type Description
DataPipeline DataPipeline

L'instance courante (pour le chaînage).

Source code in kadi/kidas/pipeline.py
def add_validation_step(
    self,
    schema: Dict[str, str],
) -> "DataPipeline":
    """Ajoute une étape de validation de schéma à la file du pipeline.

    Args:
        schema (dict[str, str]): Schéma de validation : nom_colonne → type.
            Exemple : {'culture': 'str', 'rendement_kg': 'float'}.

    Returns:
        DataPipeline: L'instance courante (pour le chaînage).
    """
    self._etapes.append({
        "type": "validation",
        "nom": "validate_schema",
        "params": {"schema": schema},
    })

    logger.debug("Étape de validation ajoutée (schéma: %s).", schema)

    return self

add_normalization_step

add_normalization_step(mappings: Dict[str, Any]) -> DataPipeline

Ajoute une étape de normalisation à la file du pipeline.

Parameters:

Name Type Description Default
mappings dict

Dictionnaire de normalisation. Les clés supportées sont 'columns' (snake_case), 'units' (unités), 'crops' (FAO), 'markets' (Bénin). Exemple : {'crops': 'culture'}.

required

Returns:

Name Type Description
DataPipeline DataPipeline

L'instance courante (pour le chaînage).

Source code in kadi/kidas/pipeline.py
def add_normalization_step(
    self,
    mappings: Dict[str, Any],
) -> "DataPipeline":
    """Ajoute une étape de normalisation à la file du pipeline.

    Args:
        mappings (dict): Dictionnaire de normalisation. Les clés supportées
            sont 'columns' (snake_case), 'units' (unités), 'crops' (FAO),
            'markets' (Bénin). Exemple : {'crops': 'culture'}.

    Returns:
        DataPipeline: L'instance courante (pour le chaînage).
    """
    self._etapes.append({
        "type": "normalization",
        "nom": "normalize",
        "params": {"mappings": mappings},
    })

    logger.debug("Étape de normalisation ajoutée (mappings: %s).", mappings)

    return self

execute

execute(cache: bool = True) -> Tuple[pd.DataFrame, Dict]

Exécute toutes les étapes configurées du pipeline.

Charge les données depuis la source, applique les étapes de nettoyage, validation et normalisation dans l'ordre, puis met le résultat en cache si demandé.

Parameters:

Name Type Description Default
cache bool

Si True, tente de charger depuis le cache avant la lecture et sauvegarde le résultat final. Par défaut True.

True

Returns:

Type Description
Tuple[DataFrame, Dict]

tuple[pd.DataFrame, dict]: Tuple contenant : - Le DataFrame traité et prêt à l'emploi. - Le rapport complet de toutes les étapes.

Raises:

Type Description
KidasPipelineError

Si aucune source n'a été configurée.

KidasReadError

Si la lecture de la source échoue.

Source code in kadi/kidas/pipeline.py
def execute(
    self,
    cache: bool = True,
) -> Tuple[pd.DataFrame, Dict]:
    """Exécute toutes les étapes configurées du pipeline.

    Charge les données depuis la source, applique les étapes de
    nettoyage, validation et normalisation dans l'ordre, puis
    met le résultat en cache si demandé.

    Args:
        cache (bool): Si True, tente de charger depuis le cache avant
            la lecture et sauvegarde le résultat final. Par défaut True.

    Returns:
        tuple[pd.DataFrame, dict]: Tuple contenant :
            - Le DataFrame traité et prêt à l'emploi.
            - Le rapport complet de toutes les étapes.

    Raises:
        KidasPipelineError: Si aucune source n'a été configurée.
        KidasReadError: Si la lecture de la source échoue.
    """
    # Vérification qu'une source a été configurée
    if self._source is None:
        raise KidasPipelineError(
            "Aucune source configurée. Appelez load_data() avant execute()."
        )

    # Génération d'une clé de cache basée sur le chemin de la source
    cle_cache = f"pipeline_{self._source.source_path}"

    # Tentative de chargement depuis le cache
    if cache:
        df_cached, _ = self._cache.load(cle_cache)
        if df_cached is not None:
            logger.info(
                "Pipeline kidas : données chargées depuis le cache (clé: '%s').",
                cle_cache,
            )
            self._rapports["cache_utilise"] = True
            return df_cached, self._rapports

    # Lecture des données depuis la source
    try:
        self._df = self._source.read()
        logger.info(
            "Pipeline kidas : %d lignes chargées depuis '%s'.",
            len(self._df),
            self._source.source_path,
        )
    except Exception as erreur:
        raise KidasReadError(
            f"Échec de lecture dans le pipeline : {erreur}"
        ) from erreur

    # Exécution des étapes dans l'ordre
    for etape in self._etapes:
        self._df = self._executer_etape(etape)

    # Sauvegarde en cache du résultat final
    if cache and self._df is not None:
        self._cache.save(cle_cache, self._df)
        self._rapports["cache_utilise"] = True

    # Compilation du rapport final
    self._rapports["lignes_finales"] = len(self._df) if self._df is not None else 0
    self._rapports["etapes_appliquees"] = [e["nom"] for e in self._etapes]

    return self._df, self._rapports

get_pipeline_config

get_pipeline_config() -> dict

Retourne la configuration complète du pipeline (étapes définies).

Returns:

Name Type Description
dict dict

Dictionnaire décrivant la source et les étapes configurées.

Source code in kadi/kidas/pipeline.py
def get_pipeline_config(self) -> dict:
    """Retourne la configuration complète du pipeline (étapes définies).

    Returns:
        dict: Dictionnaire décrivant la source et les étapes configurées.
    """
    return {
        "source": self._rapports.get("source"),
        "nb_etapes": len(self._etapes),
        "etapes": [
            {"type": e["type"], "nom": e["nom"], "params": e["params"]}
            for e in self._etapes
        ],
    }

export_report

export_report(filepath: str) -> bool

Exporte le rapport de pipeline dans un fichier JSON ou HTML.

Parameters:

Name Type Description Default
filepath str

Chemin de destination du rapport. L'extension détermine le format : '.json' ou '.html'.

required

Returns:

Name Type Description
bool bool

True si l'export s'est déroulé avec succès.

Raises:

Type Description
KidasPipelineError

Si l'extension n'est pas supportée.

Source code in kadi/kidas/pipeline.py
def export_report(self, filepath: str) -> bool:
    """Exporte le rapport de pipeline dans un fichier JSON ou HTML.

    Args:
        filepath (str): Chemin de destination du rapport. L'extension
            détermine le format : '.json' ou '.html'.

    Returns:
        bool: True si l'export s'est déroulé avec succès.

    Raises:
        KidasPipelineError: Si l'extension n'est pas supportée.
    """
    import json

    _, extension = os.path.splitext(filepath.lower())

    try:
        if extension == ".json":
            # Export au format JSON
            with open(filepath, "w", encoding="utf-8") as f:
                json.dump(self._rapports, f, ensure_ascii=False, indent=2, default=str)

        elif extension == ".html":
            # Export au format HTML simplifié
            contenu_json = json.dumps(
                self._rapports, ensure_ascii=False, indent=2, default=str
            )
            html = (
                "<html><head><meta charset='utf-8'>"
                "<title>Rapport kidas Pipeline</title></head>"
                "<body><h1>Rapport kidas DataPipeline</h1>"
                f"<pre>{contenu_json}</pre></body></html>"
            )
            with open(filepath, "w", encoding="utf-8") as f:
                f.write(html)

        else:
            raise KidasPipelineError(
                f"Format d'export '{extension}' non supporté. "
                f"Utilisez '.json' ou '.html'."
            )

        logger.info("Rapport pipeline exporté vers '%s'.", filepath)
        return True

    except OSError as erreur:
        raise KidasPipelineError(
            f"Impossible d'écrire le rapport vers '{filepath}' : {erreur}"
        ) from erreur