Data‑Lake Ingestion Pipelines in Rust – De CSV à Parquet via calamine et zip

Par Emmanuel Forgues - 6 août 2025

Chapô (96<0xE2><0x80><0xAF>mots)<0xC2><0xA0>: Face à l'explosion des volumes de données brutes, la construction d’un data‑lake fiable repose sur des pipelines d’ingestion performants. Le langage Rust, reconnu pour sa sûreté mémoire et ses performances proches du C, attire les équipes DevSecOps cherchant à allier vitesse d’exécution, faible empreinte CPU et absence de fuites ou de débordements. Cet article analyse la chaîne complète –<0xE2><0x80><0xAF>lecture d’un fichier CSV (ou XLSX) contenu dans une archive ZIP, transformation en colonnes typées, écriture au format Parquet –<0xE2><0x80><0xAF>via les crates calamine, zip et parquet‑rust. Sont détaillés l'architecture, les bonnes pratiques de gouvernance des données et les points de vigilance pour un déploiement en environnement cloud ou on‑premise.

Publié initialement le 6 août 2025.

Mis à jour le 18 avril 2026.

Migré vers StratoSentry le 1er mai 2026.

Introduction : pourquoi repenser l’ingestion ?

Les data‑lakes modernes constituent le socle analytique des organisations exploitant leurs données opérationnelles, marketing ou IoT. L’ingestion, première étape du processus, détermine la qualité, la traçabilité et les coûts d’exploitation du lac de données.

Ces pipelines sont traditionnellement écrits en Python (pandas) ou en Java/Scala (Spark). Malgré leur puissance, ils présentent souvent des limites :

  • Temps d’exécution élevés sur des volumes > 10 TB (sérialisation/désérialisation lourde).
  • Gestion de la mémoire difficile, source de pannes en production.
  • Surface d’attaque élargie : dépendances dynamiques, interpréteur Python exposant le processus à l’injection de code.

Rust offre une alternative<0xE2><0x80><0xAF>: un binaire natif, sans runtime, garantissant l’absence de data races grâce au système de possession et compilant en un exécutable compact. L'association de Rust avec les crates calamine (lecture d’Excel/CSV), zip (décompression) et parquet‑rust (écriture Parquet) permet de créer une chaîne d’ingestion capable de :

  • Traiter des archives compressées directement depuis le stockage objet (S3, Azure Blob).
  • Normaliser les schémas à la volée grâce aux métadonnées du fichier source.
  • Produire un format colonne‑orienté (Parquet) optimisé pour les requêtes analytiques et la compression.

Cet article s’adresse aux décideurs techniques (DSI, architectes data), aux équipes DevSecOps et aux développeurs évaluant Rust comme socle d’ingestion pour leurs data‑lakes.

1. Architecture générale du pipeline

+-------------------+      +-----------------+      +--------------------+
|   Source ZIP (S3) | ---> |   Rust binary   | ---> |   Parquet on Lake  |
|   (CSV / XLSX)    |      |   (calamine,    |      |   (e.g. S3/ADLS)   |
|                   |      |    zip, arrow) |      +--------------------+
+-------------------+      +-----------------+

1.1. Étapes clés

ÉtapeDescriptionCrate principale
ExtractionLecture du flux ZIP depuis le stockage objet (streaming)zip
Parsing CSV/XLSXDécodage des lignes, inférence de type (string → int/float/date)calamine
Conversion ArrowConstruction d’un RecordBatch (Apache Arrow) en mémoirearrow-rs
Écriture ParquetSérialisation du batch au format colonne‑orienté, compression Snappyparquet
Upload finalEnvoi du fichier .parquet vers le data‑lake (S3 SDK)aws-sdk-s3*

\*Le SDK AWS est utilisé à titre d’exemple<0xE2><0x80><0xAF>; la même logique s’applique avec Azure ou GCP.

1.2. Flux de données

  • Streaming : Le binaire Rust ouvre le flux S3 en mode range pour ne charger que les parties nécessaires du ZIP (optimisation réseau).
  • Décompression à la volée : zip::read::ZipArchive expose chaque entrée comme un lecteur (Read) sans écrire de fichiers temporaires sur disque.
  • Parsing : calamine::Reader lit le CSV ou les feuilles Excel et génère des vecteurs de valeurs typées, en appliquant une règle d’inférence (ex. 0/1 → bool, dates reconnues via RFC 3339).
  • Arrow Batch : Les colonnes sont transformées en ArrayRef puis assemblées dans un RecordBatch. Arrow garantit la compatibilité avec le format Parquet et facilite la compression.
  • Parquet Writer : Le writer crée un fichier .parquet local (ou en mémoire via Vec<u8>), applique la compression Snappy ou ZSTD, et ajoute les métadonnées de schéma.
  • Persistage : Le binaire utilise le SDK du cloud pour pousser le résultat final vers le bucket cible, où il sera immédiatement lisible par Athena, Presto, Spark, etc.

2. Lecture d’archives ZIP avec la crate zip

2.1. Pourquoi éviter les fichiers temporaires ?

Dans un environnement serveur sans stockage persistant (ex. conteneurs stateless), écrire chaque archive sur disque augmente le temps I/O et le risque de saturation du volume éphémère. La crate zip permet de travailler en streaming :

use std::io::{Read, Cursor};
use zip::read::ZipArchive;

// `bytes` provient d’un appel S3 GetObject (streaming)
let mut cursor = Cursor::new(bytes);
let mut archive = ZipArchive::new(&mut cursor)?; // Ouvre l’archive en lecture

for i in 0..archive.len() {
    let mut file = archive.by_index(i)?;
    if file.name().ends_with(".csv") || file.name().ends_with(".xlsx") {
        // Traitement direct du flux
        process_source(&mut file)?;
    }
}

Points de vigilance<0xE2><0x80><0xAF>:

RisqueMitigation
ZIP bomb (décompression exponentielle)Limiter le nombre d’entrées (archive.len()) et la taille totale (file.size() < MAX_BYTES).
Encodage non‑UTF8Detecter l’encodage avec encoding_rs avant parsing.

3. Parsing CSV / XLSX grâce à calamine

3.1. Fonctionnalités principales

  • Lecture native des formats Excel (xls, xlsx) et CSV (via le lecteur générique).
  • Support de la détection d’en‑têtes, conversion automatique en types Rust (i64, f64, String).
  • Gestion de cellules vides : renvoie None (option) pour permettre l’imputation ultérieure.

3.2. Exemple de code

use calamine::{open_workbook_auto, Reader, DataType};

fn parse_excel<R: Read>(mut source: R) -> Result<Vec<RecordBatch>, Box<dyn std::error::Error>> {
    let mut workbook = open_workbook_auto(&mut source)?;
    let sheet_names = workbook.sheet_names().to_owned();

let mut batches = Vec::new();
    for name in sheet_names {
        if let Some(Ok(range)) = workbook.worksheet_range(&name) {
            // Inférence de schéma
            let schema = infer_schema(&range);
            // Construction du RecordBatch Arrow
            let batch = build_record_batch(&range, &schema)?;
            batches.push(batch);
        }
    }
    Ok(batches)
}

Points d’attention<0xE2><0x80><0xAF>:

AspectDétail
Performancecalamine lit ligne par ligne ; pour des fichiers > 500 MB, envisager le streaming via csv::ReaderBuilder.
Gestion de la mémoireConvertir chaque colonne en Vec<T> puis libérer les lignes intermédiaires évite un pic d’utilisation.
Compatibilitécalamine ne supporte pas les macros Excel (.xlsm) contenant du VBA – prévoir une validation préalable.

4. Construction d’un RecordBatch Arrow

Apache Arrow est la couche de sérialisation en‑mémoire qui aligne les colonnes sur des buffers contigus, facilitant le passage à Parquet.

use arrow::array::{Int64Array, Float64Array, StringArray};
use arrow::datatypes::{DataType, Field, Schema};
use arrow::record_batch::RecordBatch;

fn build_record_batch(
    range: &calamine::Range<DataType>,
    schema: &Schema,
) -> Result<RecordBatch, Box<dyn std::error::Error>> {
    // Exemple simplifié<0xE2><0x80><0xAF>: 3 colonnes (id:int, value:float, label:string)
    let mut ids = Vec::new();
    let mut values = Vec::new();
    let mut labels = Vec::new();

for row in range.rows() {
        ids.push(row[0].get_int().unwrap_or_default());
        values.push(row[1].get_float().unwrap_or_default());
        labels.push(row[2].get_string().unwrap_or("").to_string());
    }

let id_array = Int64Array::from(ids);
    let value_array = Float64Array::from(values);
    let label_array = StringArray::from(labels);

RecordBatch::try_new(
        Arc::new(schema.clone()),
        vec![
            Arc::new(id_array) as ArrayRef,
            Arc::new(value_array),
            Arc::new(label_array),
        ],
    )
}

4.1. Inférence de schéma

L’inférence s’appuie sur les premières N lignes (ex. 1000). Si une colonne contient des entiers et des décimaux, le type est promu en Float64. Les dates au format ISO‑8601 sont converties en Timestamp(Nanosecond).

4.2. Gestion des valeurs manquantes

Arrow utilise des bitmaps de validité : chaque colonne possède un masque indiquant si les éléments sont valides (true) ou nulls (false). Les moteurs de requête ignorent ainsi les cellules vides sans pénalité de stockage.

5. Écriture au format Parquet avec parquet‑rust

5.1. Choix du codec de compression

CodecRatio moyen (texte)CPU / décompressionUsage recommandé
Snappy2–3×FaibleChargement rapide, faible latence
ZSTD4–5×ModéréStockage à long terme, volume important
GZIP> 6×ÉlevéArchivage, non‑critique en temps réel

Le writer Rust expose l’option via WriterProperties.

use parquet::file::properties::WriterProperties;
use parquet::basic::Compression;

let props = WriterProperties::builder()
    .set_compression(Compression::ZSTD)
    .build();

let mut buffer = Vec::new(); // ou un fichier temporaire
{
    let mut writer = ArrowWriter::try_new(&mut buffer, schema.clone(), Some(props))?;
    writer.write(&record_batch)?;
    writer.close()?;
}

5.2. Métadonnées et gouvernance

  • Schema evolution : Parquet supporte l’ajout de colonnes sans réécriture du fichier (via append).
  • Partitionnement : Le pipeline peut créer des dossiers /year=2024/month=03/ afin que les moteurs de requête effectuent le pruning.
  • Data‑lineage : En incluant dans les métadonnées du fichier (key_value_metadata) le nom de l’archive source, la version du schéma et le hash SHA‑256 du contenu d’origine, on assure une traçabilité complète (exigence GDPR et audit).

6. Déploiement et orchestration

6.1. Conteneurisation

Le binaire Rust se compile en un exécutable statique ; il peut être empaqueté dans une image Docker minimale (FROM scratch ou alpine). Exemple de Dockerfile :

FROM rust:1.77 as builder
WORKDIR /app
COPY . .
RUN cargo build --release

FROM gcr.io/distroless/cc
COPY --from=builder /app/target/release/ingest /usr/local/bin/ingest
ENTRYPOINT ["/usr/local/bin/ingest"]

Cette approche permet :

  • Faible surface d’attaque : aucune shell, aucun package manager.
  • Démarrage instantané (≤ 200 ms).

6.2. Orchestration via Kubernetes

Job data-ingestion :

ChampValeur
restartPolicyOnFailure
resources.limits.cpu500m (0,5 vCPU)
resources.limits.memory256Mi
volumeMountsConfigMap contenant les paramètres de connexion S3.

Le job s’exécute à chaque arrivée d’un nouvel objet dans le bucket (déclencheur S3 Event Notification → AWS Lambda qui envoie un message à une file SQS, consommée par le pod Kubernetes).

7. Cas d’usage réel : ingestion de rapports financiers mensuels

7.1. Contexte

Une société de services financiers reçoit chaque fin de mois des dizaines de fichiers ZIP contenant :

  • Un fichier report.xlsx avec les transactions (≈ 2 M lignes).
  • Un CSV metadata.csv décrivant la source et le périmètre.

Les exigences :

  • Délai : disponibilité du parquet dans le data‑lake sous 30 minutes après réception.
  • Intégrité : chaque ligne doit être horodatée avec la version du fichier source (hash SHA‑256).
  • Sécurité : aucun processus ne doit écrire sur un disque partagé.

7.2. Implémentation

  • S3 Event → SQS déclenche le job Kubernetes.
  • Le pod télécharge le ZIP en streaming (range = 0‑10 Mo) et l’ouvre via zip.
  • calamine lit la feuille « Transactions », infère les types (i64, f64, Timestamp).
  • Un RecordBatch est créé, enrichi d’une colonne source_hash.
  • Le batch est écrit en Parquet ZSTD et stocké dans s3://lake/finance/year=2024/month=03/report.parquet.
  • Une table Athena externe pointe sur le préfixe, rendant les données immédiatement interrogeables.

7.3. Résultats

KPIValeur mesurée
Temps total d’ingestion18 minutes (incluant téléchargement)
CPU utilisé< 30 % d’un vCPU pendant l’exécution
Mémoire maximale180 MiB (sans swap)
Taille du fichier Parquet620 Mo (compression ZSTD, ratio ≈ 3,5×)

Le gain de performance par rapport à un script Python/pandas (≈ 45 minutes) réduit les coûts d’infrastructure et permet le respect des SLA.

8. Points de vigilance et limites

8.1. Complexité du schéma évolutif

  • Problème : Si les colonnes changent fréquemment (ajout/suppression), le code d’inférence doit être versionné et testé.
  • Solution : Centraliser la définition de schéma dans un fichier JSON/YAML partagé entre producteurs et consommateurs; le pipeline valide chaque batch contre ce contrat.

8.2. Gestion des encodages non‑ASCII

calamine suppose UTF‑8 pour les CSV. Les jeux de données provenant d’environnements legacy (ISO‑8859‑1, Windows‑1252) nécessitent un pré‑traitement<0xE2><0x80><0xAF>:

let decoded = encoding_rs::WINDOWS_1252.decode_without_bom_handling(&raw_bytes);

Sans cette étape, les caractères spéciaux seront corrompus, entraînant des erreurs de parsing et des pertes d’information.

8.3. Sécurité du traitement ZIP

  • ZIP bomb : Un fichier contenant un petit nombre d’entrées mais une taille décompressée astronomique (ex. 1 Go → 100 TiB).
  • Mesure : Limiter file.uncompressed_size() et abortir si dépassement du seuil (ex. 5 GiB).

8.4. Compatibilité des types Arrow ↔ Parquet

Certaines variantes de type (Decimal128, Interval) ne sont pas encore pleinement supportées par le crate parquet‑rust. Dans ces cas, il faut soit :

  • Convertir en

Retour au blog

StratoSentry - 125 boulevard Saint-Denis, 92400 Courbevoie, France - contact@stratosentry.com