Aller au contenu principal
En savoir plus →
Pipeline de données

Le pipeline ETL

Chaque jeu de données d'OpenCongoData entre par un seul pipeline déclaratif et auditable — d'un simple manifeste JSON jusqu'à une table versionnée et interrogeable dans Postgres. Aucun scraper sur mesure, aucun téléversement manuel.

Ajoutez une source de données sans toucher au code

4

étapes du pipeline

6

outils dans la chaîne

1

fichier JSON par source

0

lignes de code pour en ajouter une

Comment les données circulent

Quatre étapes entièrement automatisées. Une source est déclarée une fois, puis ingérée selon un calendrier sans aucune intervention.

  1. 01

    Déclarer

    Une source et ses jeux de données sont décrits dans un seul manifeste JSON, versionné dans le dépôt.

  2. 02

    Extraire

    Un worker planifié détecte chaque manifeste et lit les données et métadonnées déclarées.

  3. 03

    Valider & transformer

    Le manifeste est validé avec Zod ; les colonnes de chaque jeu sont vérifiées via un Table Schema Frictionless puis normalisées.

  4. 04

    Charger

    Les métadonnées sont insérées dans Postgres et une nouvelle version immuable du jeu de données est ajoutée.

La chaîne d'outils

Chaque brique est open source et s'exécute dans un unique service worker Node.js.

TypeScript

Sûreté de typage de bout en bout, du schéma du manifeste jusqu'aux lignes en base.

BullMQ

File de jobs adossée à Redis, avec tâches répétables (cron), reprises et contrôle de concurrence.

Redis · ioredis

Magasin de la file et de l'état du planificateur, accédé via ioredis.

Zod

Validation à l'exécution de chaque manifeste avant la moindre écriture.

Frictionless Table Schema

Un contrat déclaratif pour les colonnes, types et contraintes de chaque jeu de données.

Supabase · PostgreSQL

Postgres avec PostGIS, recherche plein texte et Row-Level Security comme système de référence.

Étape 1 — Déclarer

Un manifeste par source

Un manifeste de source est un simple fichier JSON dans apps/workers/src/sources/. Il nomme l'éditeur et sa licence, puis liste un ou plusieurs jeux de données — chacun avec des métadonnées bilingues, un niveau d'accès, une fréquence de mise à jour (une durée ISO 8601), l'emplacement des données et un Table Schema décrivant chaque colonne. Ajouter une nouvelle source se résume donc à un fichier et une pull request — jamais à une modification du code du pipeline.

apps/workers/src/sources/health-facilities.json
{
  "source_id": "health-facilities-2024",
  "name": "DRC Health Facilities 2024",
  "publisher": "OpenCongoData Team",
  "license": "CC-BY-4.0",
  "datasets": [
    {
      "slug": "health-facilities-drc-2024",
      "title_en": "DRC Health Facilities 2024",
      "title_fr": "Établissements de santé en RDC 2024",
      "sector": "health",
      "access_level": "public",
      "update_frequency": "P1Y",
      "format": "csv",
      "license": "CC-BY-4.0",
      "file_url": "https://congodata.app/data/health-facilities-drc-2024.csv",
      "tableSchema": {
        "fields": [
          { "name": "id", "type": "string", "constraints": { "required": true, "unique": true } },
          { "name": "name", "type": "string", "constraints": { "required": true } },
          { "name": "province", "type": "string", "constraints": { "required": true } },
          { "name": "latitude", "type": "number" },
          { "name": "longitude", "type": "number" },
          { "name": "operational", "type": "boolean" }
        ],
        "primaryKey": "id"
      }
    }
  ]
}

Ajouter une source de données = un fichier JSON + une pull request.

Étape 2 — Orchestrer

Planifié, repris, sans chevauchement

Au démarrage, le worker détecte chaque manifeste et enregistre pour lui un job répétable BullMQ — par défaut, chaque jour à 02h00. Les jobs s'exécutent sur une file adossée à Redis : ils survivent aux redémarrages, sont automatiquement repris en cas d'échec et ne se chevauchent jamais. Le même job peut aussi être lancé à la demande pour une actualisation immédiate.

apps/workers/src/queue.ts — la file & le worker
import { Queue, Worker } from "bullmq";
import { Redis } from "ioredis";

const connection = new Redis(
  process.env.REDIS_URL ?? "redis://localhost:6379",
  { maxRetriesPerRequest: null },
);

export const ingestionQueue = new Queue("ingestion", { connection });

export function createIngestionWorker(
  processor: (job: { data: { sourceId: string } }) => Promise<void>,
) {
  return new Worker("ingestion", (job) => processor(job), { connection });
}
apps/workers/src/index.ts — planifier chaque source
// Discover every manifest and register a repeatable job.
for (const sourceId of sourceIds) {
  await ingestionQueue.add(
    "ingest",
    { sourceId },
    { repeat: { pattern: "0 2 * * *" } }, // daily at 02:00
  );
}

const worker = createIngestionWorker(async (job) => {
  await ingestSource(job.data.sourceId);
});

worker.on("failed", (job, err) => {
  console.error(`Job ${job?.id} failed:`, err);
});

Les échecs sont remontés via l'événement failed du worker, prêts pour le journal et les alertes.

Étape 3 — Extraire & valider

Rien n'est fiable tant que ce n'est pas validé

Pour chaque job, le worker lit le manifeste, l'analyse et le passe par SourceManifestSchema — un schéma Zod qui fait autorité sur la forme d'une source. Un manifeste mal formé échoue ici, bruyamment, avant toute écriture en base. Chaque jeu de données porte en outre un Table Schema Frictionless : le contrat servant à valider et normaliser les données tabulaires sous-jacentes, jusqu'aux noms de colonnes, aux types et aux contraintes required / unique.

apps/workers/src/ingest.ts — lire, analyser, valider
import { SourceManifestSchema } from "@opencongodata/schemas";
import { readFile } from "node:fs/promises";

export async function ingestSource(sourceId: string) {
  const path = resolve(__dirname, "sources", `${sourceId}.json`);
  const raw = await readFile(path, "utf-8");

  // A malformed manifest throws here — before any DB write.
  const manifest = SourceManifestSchema.parse(JSON.parse(raw));
  // … extract + load follow
}
packages/schemas/src/dataset.ts — le contrat du manifeste
// Zod is the single source of truth for a manifest's shape.
export const DataFormatSchema = z.enum([
  "csv", "json", "geojson", "parquet", "xlsx",
]);

export const DatasetManifestSchema = z.object({
  slug:            z.string().regex(/^[a-z0-9-]+$/),
  title_fr:        z.string().min(1),
  title_en:        z.string().min(1),
  sector:          SectorSchema,
  access_level:    AccessLevelSchema,
  update_frequency: z.string().regex(/^P/),  // ISO 8601 duration
  file_url:        z.string().url().optional(),
  format:          DataFormatSchema,
  license:         z.string().min(1),
  tableSchema:     TableSchemaSchema,        // Frictionless contract
});

Les champs du Table Schema acceptent les types string, number, integer, boolean, date, datetime et geojson — ce dernier débloquant la géométrie PostGIS.

Étape 4 — Charger

Mettre à jour les métadonnées, ajouter une version

Les données validées sont écrites avec le client Supabase à rôle de service. La source et chaque jeu de données sont mis à jour (upsert) — indexés par source_id et slug — de sorte que relancer l'ingestion est sûr et idempotent. Chaque exécution ajoute ensuite une ligne à dataset_versions, donnant à chaque jeu un historique immuable et horodaté.

apps/workers/src/ingest.ts — charger dans Postgres
const db = getDb(); // service-role Supabase client

// 1. Upsert the source (idempotent, keyed by source_id).
await db.from("sources").upsert(
  {
    source_id: manifest.source_id,
    name:      manifest.name,
    publisher: manifest.publisher,
    license:   manifest.license,
  },
  { onConflict: "source_id" },
);

for (const dataset of manifest.datasets) {
  // 2. Upsert dataset metadata (keyed by slug).
  const { data: ds } = await db.from("datasets").upsert(
    {
      slug:      dataset.slug,
      sector:    dataset.sector,
      format:    dataset.format,
      source_id: manifest.source_id,
      // …bilingual title/description, access_level, license
    },
    { onConflict: "slug" },
  ).select("id").single();

  // 3. Append an immutable version record.
  await db.from("dataset_versions").insert({
    dataset_id:   ds!.id,
    version:      Date.now(),
    storage_path: dataset.file_url ?? null,
  });
}

Idempotent

Les upserts indexés sur des identifiants stables mettent à jour sur place plutôt que de dupliquer.

Versionné

Chaque ingestion ajoute à dataset_versions — l'historique est conservé, jamais écrasé.

Côté serveur uniquement

La clé de rôle de service ne vit que dans le worker ; chaque accès client passe par la Row-Level Security.

Exécuter en local

Le worker est un service Node.js standard. Pointez-le vers une instance Redis et votre projet Supabase, puis lancez-le seul ou avec le reste de la stack.

.env — variables requises
REDIS_URL=redis://localhost:6379
SUPABASE_URL=http://localhost:54321
SUPABASE_SERVICE_ROLE_KEY=<your-service-role-key>
Démarrer le worker
# Run just the ingestion worker
pnpm --filter @opencongodata/workers dev

# …or the whole stack (web + api + workers)
pnpm dev

Ajouter votre propre source

Le pipeline étant déclaratif, contribuer des données ne demande jamais de toucher au code d'ingestion. En trois étapes :

  1. 1

    Créez un manifeste JSON dans apps/workers/src/sources/ décrivant votre source et ses jeux de données.

  2. 2

    Validez-le en local — le schéma Zod rejette tout ce qui est mal formé avant publication.

  3. 3

    Ouvrez une pull request. Une fois fusionnée, la source est planifiée et ingérée automatiquement.