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.