Retour aux Guides

    É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.

    client.producer()client.consumer("group")client.admin()

    1. Démarrage rapide

    Démarrez un cluster local, puis envoyez un événement. Dix lignes. C'est tout.

    bash
    1# Start a local StreamFlow cluster
    2./scripts/start-streamflow.sh
    3# StreamFlow is listening on localhost:9092
    Java
    1StreamFlowClient 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.

    Développer localement
    StreamFlowClient.embedded(...)
    Application
        ↓
    StreamFlow SDK
        ↓
    Engine (in-process, ~µs)

    Tests, benchmarks, monolithes.

    Déployer
    StreamFlowClient.remote(...)
    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

    Java
    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 ClientStreamFlow 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

    XML
    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

    Java
    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

    Java
    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

    Java
    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

    Java
    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

    Java
    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.

    Java
    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.

    Java
    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.

    bash
    1# Start a cluster (or point at an existing one)
    2./scripts/start-streamflow.sh
    3# Listening on port 9092
    Java
    1StreamFlowClient 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.

    Java
    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é

    Java
    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

    Java
    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

    Java
    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

    Java
    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

    Java
    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/