Écrire une fois.
Exécuter embedded.
Passer à l'échelle en remote.
Le même code tourne en embedded pour les tests et les benchmarks, puis se connecte à un cluster StreamFlow distant en production. Votre code métier reste identique, seules la construction du client et quelques appels de cycle de vie propres au transport diffèrent.
Sur cette page
1. Démarrage rapide
Démarrez un cluster local, puis envoyez un événement. Dix lignes. C'est tout.
1# Start a local StreamFlow cluster
2./scripts/start-streamflow.sh
3# StreamFlow is listening on localhost:90921StreamFlowClient client = StreamFlowClient.remote("localhost:9092");
2
3try (var producer = client.producer()) {
4 producer.send(
5 "orders",
6 "42",
7 orderBytes
8 ).join();
9}L'événement a été acquitté par StreamFlow et est disponible pour tous les consumers du topic. La durabilité dépend du profil de stockage et de réplication actif sur le cluster.
2. Développer localement, déployer à distance
Une API. Deux transports. Aucune réécriture.
Application
↓
StreamFlow SDK
↓
Engine (in-process, ~µs)Tests, benchmarks, monolithes.
Application
↓
StreamFlow SDK
↓
TCP (Kafka wire)
↓
Cluster (~ms)Microservices, multi-JVM, production.
La promesse : votre code métier reste identique. Seules la construction du client et quelques appels de cycle de vie propres au transport (flush, commitSync, groupes de consumers, propriétés personnalisées) diffèrent entre les deux modes.
Même code métier, deux transports
1StreamFlowClient client =
2 environment.equals("prod")
3 ? StreamFlowClient.remote(System.getenv("STREAMFLOW_BOOTSTRAP"))
4 : StreamFlowClient.embedded(router, router, eventLog);
5
6OrderService service = new OrderService(client);
7service.submit("order-42", orderBytes);3. Interopérabilité Kafka
Vous venez de Kafka ? Chaque méthode a un équivalent direct.
| Apache Kafka Client | StreamFlow SDK |
|---|---|
| new KafkaProducer<>(props) | client.producer() |
| producer.send(new ProducerRecord<>(...)) | producer.send("topic", key, val) |
| new KafkaConsumer<>(props) | client.consumer("group") |
| consumer.subscribe(List.of("topic")) | consumer.subscribe("topic") |
| consumer.poll(Duration) | consumer.poll(Duration) |
| consumer.commitSync() | consumer.commitSync() |
| AdminClient.create(props) | client.admin() |
StreamFlow implémente le protocole filaire Kafka. N'importe quel client Kafka en Java, Python, Go ou Node.js peut lire ce que le SDK écrit et réciproquement.
4. Installation
1<properties>
2 <streamflow.version>0.1.0</streamflow.version>
3</properties>
4
5<dependency>
6 <groupId>com.streamflow</groupId>
7 <artifactId>streamflow-sdk</artifactId>
8 <version>${streamflow.version}</version>
9</dependency>- Remote nécessite aussi
kafka-clients - Embedded nécessite aussi
streamflow-core - Les deux sont déclarés
<optional>dans le SDK, n'ajoutez que ce dont vous avez besoin.
5. Produire des événements
Envoyer un événement
1try (var producer = client.producer()) {
2 RecordMetadata meta = producer
3 .send("orders", "order-42", orderBytes)
4 .join();
5
6 System.out.printf("partition=%d offset=%d%n",
7 meta.partition(), meta.offset());
8}Envoyer avec des headers
1producer.send(
2 "orders",
3 "order-42",
4 orderBytes,
5 Map.of(
6 "source", "mobile-app",
7 "region", "eu-west",
8 "trace-id", "abc-123"
9 )
10).join();Avancé : boucle par lots
1try (var producer = client.producer()) {
2 for (int i = 1; i <= 10_000; i++) {
3 producer.send("orders", "order-" + i, payload(i));
4 }
5 producer.flush(); // remote only, force network send
6}6. Consommer des événements
Recevoir un lot
1try (var consumer = client.consumer("order-service")) {
2 consumer.subscribe("orders");
3
4 List<StreamFlowRecord> records =
5 consumer.poll(Duration.ofSeconds(1));
6
7 for (var r : records) {
8 process(r.key(), r.value());
9 }
10 consumer.commitSync();
11}Avancé : boucle continue
1var consumer = client.consumer("payment-processor");
2consumer.subscribe("payments");
3
4while (running) {
5 var records = consumer.poll(Duration.ofMillis(100));
6 for (var r : records) {
7 handle(r);
8 }
9 if (!records.isEmpty()) consumer.commitSync();
10}7. Rejeu
Revenir à un offset donné sur une partition et rejouer.
1consumer.subscribe("orders");
2consumer.seek("orders", 0, 500L);
3
4var replayed = consumer.poll(Duration.ofSeconds(5));8. Mode Embedded
StreamFlow tourne dans la même JVM que votre application. Idéal pour tests, benchmarks et monolithes.
1StreamFlowApplication app = new StreamFlowApplication(
2 "streamflow.properties");
3app.start();
4
5StreamFlowClient client = StreamFlowClient.embedded(
6 app.getRouter(),
7 app.getRouter(),
8 app.getEventLog()
9);
10
11try (var admin = client.admin()) {
12 admin.createTopic("temperature-readings", 4);
13}
14
15try (var producer = client.producer()) {
16 producer.send("temperature-readings", "sensor-A", reading);
17 // no flush() needed, writes are synchronous
18}
19
20app.close();Différences avec remote : pas de flush(), pas de consumer group requis, pas de commitSync() requis. Aucun saut réseau. Les écritures sont synchrones et l'exécution est déterministe. Le stockage suit la configuration du moteur embedded (mémoire ou journal local).
9. Mode Remote
Votre application se connecte à un cluster StreamFlow via TCP en utilisant le protocole filaire Kafka.
1# Start a cluster (or point at an existing one)
2./scripts/start-streamflow.sh
3# Listening on port 90921StreamFlowClient client = StreamFlowClient.remote(
2 "host1:9092,host2:9092");
3
4try (var producer = client.producer()) {
5 producer.send("orders", "42", orderBytes).join();
6 producer.flush();
7}10. Patterns
Pattern 1 : Service agnostique du mode
Logique métier écrite une fois. Injectée en embedded pour les tests, en remote en production.
1public class OrderService {
2 private final StreamFlowProducer producer;
3 private final StreamFlowConsumer consumer;
4
5 public OrderService(StreamFlowClient client) {
6 this.producer = client.producer();
7 this.consumer = client.consumer("order-service");
8 this.consumer.subscribe("orders");
9 }
10
11 public CompletableFuture<RecordMetadata> submit(String id, byte[] data) {
12 return producer.send("orders", id, data);
13 }
14}
15
16// Test
17var testClient = StreamFlowClient.embedded(router, router, eventLog);
18new OrderService(testClient);
19
20// Production
21var prodClient = StreamFlowClient.remote("prod-cluster:9092");
22new OrderService(prodClient);Pattern 2 : Rejeu depuis un offset donné
1var consumer = client.consumer("replay-group");
2consumer.subscribe("orders");
3consumer.seek("orders", 0, 500L);
4
5var records = consumer.poll(Duration.ofSeconds(5));
6consumer.commitSync();11. Référence API
StreamFlowClient
1StreamFlowClient.remote("host1:9092,host2:9092");
2StreamFlowClient.embedded(publisher, reader, eventLog);
3
4client.producer(); // default config
5client.producer(props); // custom config (remote only)
6client.consumer("group-id");
7client.consumer("group-id", props);
8client.admin();StreamFlowProducer
1CompletableFuture<RecordMetadata> f = producer.send("topic", "key", value);
2producer.send("topic", "key", value, Map.of("trace-id", "abc"));
3RecordMetadata meta = producer.send("topic", "key", value).join();
4producer.flush();
5producer.close();StreamFlowConsumer
1consumer.subscribe("orders");
2consumer.subscribe("inventory", "shipping");
3
4List<StreamFlowRecord> records = consumer.poll(Duration.ofMillis(100));
5
6for (StreamFlowRecord r : records) {
7 r.topic(); r.partition(); r.offset();
8 r.key(); r.value(); r.headers(); r.timestamp();
9}
10
11consumer.commitSync();
12consumer.seek("orders", 0, 1000L);
13long latest = consumer.latestOffset("orders", 0);
14long earliest = consumer.earliestOffset("orders", 0);StreamFlowAdmin
1try (StreamFlowAdmin admin = client.admin()) {
2 admin.createTopic("orders", 8);
3 boolean exists = admin.topicExists("orders");
4 TopicInfo info = admin.describeTopic("orders");
5 admin.deleteTopic("orders");
6}Dépendances de modules
streamflow-sdk ├── streamflow-common (required) , StreamEvent, EventPublisher, EventReader ├── streamflow-core (optional) , ShardRouter, EventLog (embedded only) └── kafka-clients (optional) , KafkaProducer/Consumer (remote only)
Prochaines étapes
- 📊 Tuning: profils de performance (haut débit, faible latence, durabilité)
- 🔌 Source Connectors: créer des sources de données personnalisées avec le SDK Connector
- 🧪 Examples: benchmarks exécutables dans examples/kafka-replication-demo/