KafkaDataDevStreaming

Kafka Streams : traitement de données en temps réel sans infrastructure séparée

9 octobre 2026 · Sphinx-Digital

Kafka Streams est une bibliothèque Java (et Scala) qui permet de traiter des flux de données Kafka sans infrastructure séparée. Pas de cluster Spark, pas de Flink, pas de Samza — juste une dépendance dans votre JAR.

KStream et KTable : les deux abstractions fondamentales

KStream représente un flux de records individuels, interprétés comme un journal d’événements. Chaque record est indépendant.

KTable représente une vue matérialisée, où chaque clé a une valeur courante. Les nouveaux records mettent à jour (ou suppriment) la valeur de la clé.

StreamsBuilder builder = new StreamsBuilder();

// KStream : chaque message est un événement
KStream<String, OrderEvent> orders = builder.stream("orders");

// KTable : vue matérialisée de l'état courant des stocks
KTable<String, StockLevel> stocks = builder.table("stock-levels");

Opérations de base sur un KStream

StreamsBuilder builder = new StreamsBuilder();

KStream<String, String> rawLogs = builder.stream("raw-logs");

// Filter : garder uniquement les erreurs
KStream<String, String> errors = rawLogs
    .filter((key, value) -> value.contains("ERROR"));

// Map : transformer chaque record
KStream<String, LogEvent> parsedErrors = errors
    .mapValues(value -> LogEvent.parse(value));

// Branch : router vers différents topics selon une condition
Map<String, KStream<String, LogEvent>> branches = parsedErrors
    .split(Named.as("severity-"))
    .branch((k, v) -> v.getSeverity() == Severity.CRITICAL,
            Branched.as("critical"))
    .branch((k, v) -> v.getSeverity() == Severity.HIGH,
            Branched.as("high"))
    .defaultBranch(Branched.as("other"));

// Envoyer vers différents topics
branches.get("severity-critical").to("critical-alerts");
branches.get("severity-high").to("high-alerts");

Fenêtres de temps : agréger sur une période

KStream<String, OrderEvent> orders = builder.stream("orders");

// Compter les commandes par heure par client
KTable<Windowed<String>, Long> orderCountPerHour = orders
    .groupByKey()
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofHours(1)))
    .count(Materialized.as("order-count-store"));

// Détecter les clients qui commandent plus de 10 fois en 1h
orderCountPerHour
    .toStream()
    .filter((windowedKey, count) -> count > 10)
    .map((windowedKey, count) -> KeyValue.pair(
        windowedKey.key(),
        new FraudAlert(windowedKey.key(), count, windowedKey.window())
    ))
    .to("fraud-alerts");

Jointure KStream × KTable

KStream<String, OrderEvent> orders = builder.stream("orders");
KTable<String, UserProfile> users = builder.table("user-profiles");

// Enrichir chaque commande avec le profil utilisateur
KStream<String, EnrichedOrder> enrichedOrders = orders.join(
    users,
    (order, user) -> new EnrichedOrder(order, user.getName(), user.getTier()),
    Joined.with(Serdes.String(), orderSerde, userProfileSerde)
);

enrichedOrders.to("enriched-orders");

Démarrer une application Streams

Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-processor");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());

StreamsBuilder builder = new StreamsBuilder();
// ... définir la topologie ...

KafkaStreams streams = new KafkaStreams(builder.build(), props);

// Gestion propre des arrêts
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));

streams.start();

Notre formation Kafka couvre Kafka Streams avec des ateliers sur des cas de traitement temps réel réels.