Migrer de Apache Kafka vers StreamFlow

    Table des matières

    1. Pourquoi migrer ?
    2. Compatibilité
    3. Plan de migration
    4. Migration par langage
    5. Migration Kafka Connect → StreamFlow Connectors
    6. Kafka Streams → Agent Graphs
    7. Checklist de migration
    8. FAQ

    1. Pourquoi migrer ?

    AspectApache KafkaStreamFlow
    Consumer throughput462K msg/s743K msg/s (+61%)
    Producer throughput447K msg/s358K msg/s (80%, gap = pd-ssd)
    Engine natif (SDK embedded)N/A3.8M evt/s
    Agents IA intégrésNon (besoin Kafka Streams + service externe)Oui (LLM, rules, MCP tools)
    Graphes d'agentsNonOui (DSL fluent, routing conditionnel)
    Elastic partitioningNon (partitions fixes)Oui (16384 slots, rebalance auto)
    Géo-réplicationMirrorMaker (outil séparé)Intégrée (active-passive + active-active)
    io_uringNon (mmap + page cache)Oui (SQPOLL, DMA buffers)
    DéploiementZooKeeper ou KRaft (complexe)Un seul binaire (Raft intégré)

    2. Compatibilité

    StreamFlow implémente le protocole Kafka wire (TCP, port 9092). Vos applications existantes se connectent sans changement de code :

    ┌──────────────────────┐        ┌──────────────────────┐
    │  Votre application   │        │  Votre application   │
    │  KafkaProducer       │        │  KafkaProducer       │  ← MÊME CODE
    │  KafkaConsumer       │        │  KafkaConsumer       │
    │  bootstrap=kafka:9092│        │  bootstrap=sf:9092   │  ← JUSTE CHANGER L'ADRESSE
    └──────────┬───────────┘        └──────────┬───────────┘
               │                               │
         Apache Kafka                    StreamFlow

    Ce qui fonctionne sans changement

    Feature KafkaSupportéNotes
    KafkaProducer (Java, Python, Go, Node.js)Toutes les libs Kafka
    KafkaConsumer with consumer groupsRebalancing, offset commit
    acks=0, acks=1, acks=allMêmes sémantiques
    batch.size, linger.msMêmes configs
    compression.type (snappy, lz4, zstd)Décompression transparente
    Topic auto-creationConfigurable
    kafka-console-producer.shOutils CLI standard
    kafka-console-consumer.sh
    kafka-producer-perf-test.sh
    kafka-consumer-perf-test.sh
    Idempotent producerInitProducerId supporté
    SASL PLAIN authentication

    Ce qui n'est pas (encore) supporté

    Feature KafkaStatusAlternative StreamFlow
    Kafka Transactionsacks=all + Raft replication
    Kafka StreamsAgent graphs (plus puissant)
    Kafka ConnectSource/Sink Connector SDK (19 built-in)
    Schema RegistryHeaders + validation SDK
    Exactly-once semantics (EOS)PartielIdempotent producer + at-least-once consumer
    ACL via kafka-acls.shREST API + TopicAclPolicy
    Log compactionLogCompactor intégré

    3. Plan de migration

    Phase 1 : Test en parallèle (1 jour)

    Faites tourner StreamFlow à côté de Kafka, sans toucher à la production.

    bash
    1# 1. Start StreamFlow (single node for testing)
    2./scripts/start-streamflow.sh
    3
    4# 2. Create the same topic as in Kafka
    5kafka-topics.sh --create --topic orders \
    6  --partitions 16 --bootstrap-server localhost:9092
    7
    8# 3. Run the same benchmark as Kafka
    9kafka-producer-perf-test.sh --topic orders \
    10  --num-records 1000000 --record-size 256 \
    11  --throughput -1 \
    12  --producer-props bootstrap.servers=localhost:9092 \
    13  batch.size=65536 linger.ms=5 acks=1
    14
    15kafka-consumer-perf-test.sh --topic orders \
    16  --messages 1000000 --broker-list localhost:9092 --timeout 60000
    17
    18# 4. Compare results with Kafka

    Phase 2 : Dual-write (1-2 semaines)

    Vos producers écrivent dans les deux clusters. Les consumers lisent de StreamFlow.

    Java
    1// Dual-write producer (during migration)
    2public class DualWriteProducer {
    3    private final KafkaProducer<String, byte[]> kafkaProducer;   // old
    4    private final StreamFlowProducer streamFlowProducer;          // new
    5
    6    public void send(String topic, String key, byte[] value) {
    7        // Write to both
    8        kafkaProducer.send(new ProducerRecord<>(topic, key, value));
    9        streamFlowProducer.send(topic, key, value);
    10    }
    11}

    Vérifier que les données sont identiques :

    bash
    1# Count events in each system
    2kafka-run-class.sh kafka.tools.GetOffsetShell \
    3  --broker-list kafka:9092 --topic orders
    4# vs
    5kafka-run-class.sh kafka.tools.GetOffsetShell \
    6  --broker-list streamflow:9092 --topic orders

    Phase 3 : Basculement (1 jour)

    bash
    1# 1. Stop Kafka consumers
    2# 2. Change bootstrap.servers in your app config
    3sed -i 's/kafka-cluster:9092/streamflow-cluster:9092/' application.yml
    4
    5# 3. Redeploy apps
    6kubectl rollout restart deployment/my-app
    7
    8# 4. Verify everything works
    9curl http://streamflow:8080/metrics | grep events_published

    Phase 4 : Décommissionner Kafka (après 1 mois)

    bash
    1# Once StreamFlow is stable in production
    2docker compose -f kafka-cluster.yml down
    3# or
    4kubectl delete statefulset kafka

    4. Migration par langage

    Java (KafkaProducer/KafkaConsumer)

    Avant (Kafka) :

    Java
    1Properties props = new Properties();
    2props.put("bootstrap.servers", "kafka-cluster:9092");
    3props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    4props.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer");
    5KafkaProducer<String, byte[]> producer = new KafkaProducer<>(props);
    6producer.send(new ProducerRecord<>("orders", "key-1", payload));

    Après (StreamFlow SDK) :

    Java
    1StreamFlowClient client = StreamFlowClient.remote("streamflow-cluster:9092");
    2StreamFlowProducer producer = client.producer();
    3producer.send("orders", "key-1", payload);

    Ou sans changer le code du tout :

    Java
    1// ONLY CHANGE: the cluster address
    2props.put("bootstrap.servers", "streamflow-cluster:9092");
    3// The rest of the KafkaProducer code works as-is

    Python (confluent-kafka)

    Avant :

    python
    1producer = Producer({'bootstrap.servers': 'kafka-cluster:9092'})

    Après :

    python
    1producer = Producer({'bootstrap.servers': 'streamflow-cluster:9092'})
    2# That's it. Same code, same lib.

    Node.js (KafkaJS)

    Avant :

    javascript
    1const kafka = new Kafka({ brokers: ['kafka-cluster:9092'] });

    Après :

    javascript
    1const kafka = new Kafka({ brokers: ['streamflow-cluster:9092'] });

    Go (kafka-go)

    Avant :

    go
    1writer := &kafka.Writer{Addr: kafka.TCP("kafka-cluster:9092"), Topic: "orders"}

    Après :

    go
    1writer := &kafka.Writer{Addr: kafka.TCP("streamflow-cluster:9092"), Topic: "orders"}

    5. Migration Kafka Connect → StreamFlow Connectors

    Kafka Connect (avant)

    JSON
    1{
    2  "name": "jdbc-source",
    3  "config": {
    4    "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
    5    "connection.url": "jdbc:postgresql://db:5432/shop",
    6    "table.whitelist": "orders",
    7    "mode": "incrementing",
    8    "incrementing.column.name": "id",
    9    "topic.prefix": "cdc-"
    10  }
    11}

    StreamFlow Connector (après)

    Properties
    1# connect-source.properties
    2source.jdbc-orders.type=jdbc
    3source.jdbc-orders.topic=cdc-orders
    4source.jdbc-orders.connection.url=jdbc:postgresql://db:5432/shop
    5source.jdbc-orders.table=orders
    6source.jdbc-orders.mode=incrementing
    7source.jdbc-orders.incrementing.column=id
    8source.jdbc-orders.poll.interval.ms=1000

    Mapping des connecteurs

    Kafka ConnectStreamFlowConfig type
    JdbcSourceConnectorJdbcSourceConnectorjdbc
    JdbcSinkConnectorJdbcSinkConnectorjdbc
    FileStreamSourceConnectorHttpSourceConnectorhttp
    ElasticsearchSinkConnectorElasticsearchSinkConnectorelasticsearch
    MongoSinkConnectorMongoDbSinkConnectormongodb
    S3SinkConnector❌ (not yet)
    DebeziumSourceConnectorCdcSourceConnectorcdc

    6. Kafka Streams → Agent Graphs

    Kafka Streams (avant)

    Java
    1StreamsBuilder builder = new StreamsBuilder();
    2
    3builder.stream("transactions")
    4    .filter((key, value) -> value.getAmount() > 10000)
    5    .mapValues(tx -> new FraudCheck(tx))
    6    .to("fraud-alerts");
    7
    8KafkaStreams streams = new KafkaStreams(builder.build(), props);
    9streams.start();

    StreamFlow Agent Graph (après)

    Java
    1CompiledGraph graph = StreamflowGraph.define("fraud-detection")
    2    .start("filter")
    3        .agent("amount-filter")
    4        .prompt("Filter: keep only transactions with amount > 10000. " +
    5               "Return the transaction unchanged or null to drop.")
    6    .then("analyze")
    7        .agent("fraud-analyzer")
    8        .prompt("Analyze this high-value transaction for fraud indicators. " +
    9               "Return {fraud: true/false, confidence: 0-1, reason: '...'}")
    10    .end("alert")
    11        .agent("alert-generator")
    12        .prompt("Generate a fraud alert for the security team. " +
    13               "Include transaction details and risk assessment.")
    14    .compile();
    15
    16GraphRunner runner = new GraphRunner();
    17runner.register(graph);
    Avantage StreamFlow : Les agents utilisent du raisonnement IA (un modèle de langage) au lieu de code impératif. Un filter Kafka Streams fait if amount > 10000. Un agent StreamFlow peut analyser le contexte, l'historique du client, les patterns de fraude, en langage naturel.

    7. Checklist de migration

    Phase 1 : Test

    • Installer StreamFlow (1 nœud)
    • Lancer le benchmark avec kafka-producer-perf-test
    • Comparer les résultats avec Kafka
    • Tester les consumer groups

    Phase 2 : Dual-write

    • Configurer le dual-write dans les producers
    • Basculer les consumers vers StreamFlow
    • Vérifier la parité des données (offsets, comptes)
    • Monitorer pendant 1-2 semaines

    Phase 3 : Basculement

    • Changer bootstrap.servers dans toutes les apps
    • Redéployer
    • Vérifier /metrics (throughput, latency, lag)
    • Configurer les alertes Prometheus

    Phase 4 : Décommissionner Kafka

    • Attendre 1 mois de stabilité
    • Exporter les offsets si nécessaire
    • Arrêter le cluster Kafka
    • Supprimer les volumes ZooKeeper/KRaft

    Bonus : Activer les features StreamFlow

    • Créer des agents IA (Guide des Agents)
    • Configurer la géo-réplication
    • Activer l'elastic partitioning
    • Migrer les Kafka Connect vers les connecteurs StreamFlow
    • Remplacer Kafka Streams par des Agent Graphs

    8. FAQ

    Q: Est-ce que je dois changer mon code ?

    R: Non. Changez juste bootstrap.servers. Toutes les libs Kafka (Java, Python, Go, Node.js) fonctionnent.

    Q: Et mes offsets Kafka ?

    R: Les offsets StreamFlow repartent de 0. Utilisez auto.offset.reset=earliest pour relire depuis le début, ou latest pour ne lire que les nouveaux messages.

    Q: Est-ce que mes consumer groups fonctionnent ?

    R: Oui. Le rebalancing, offset commit, et session timeout fonctionnent comme dans Kafka.

    Q: Performance comparée ?

    R: Consumer 61% plus rapide que Kafka. Producer à 80% de Kafka sur pd-ssd (gap dû à la latence réseau du disque, pas au code). Sur NVMe local, StreamFlow devrait égaler ou dépasser Kafka.

    Q: Et la durabilité ?

    R: Avec acks=all, les données sont répliquées via Raft sur 3 nœuds avant ACK. Même garantie que Kafka avec replication.factor=3.

    Q: Puis-je revenir à Kafka ?

    R: Oui. Pendant la phase dual-write, les deux clusters ont les mêmes données. Il suffit de rechanger bootstrap.servers.