Migrer de Apache Kafka vers StreamFlow
Table des matières
- Pourquoi migrer ?
- Compatibilité
- Plan de migration
- Migration par langage
- Migration Kafka Connect → StreamFlow Connectors
- Kafka Streams → Agent Graphs
- Checklist de migration
- FAQ
1. Pourquoi migrer ?
| Aspect | Apache Kafka | StreamFlow |
|---|---|---|
| Consumer throughput | 462K msg/s | 743K msg/s (+61%) |
| Producer throughput | 447K msg/s | 358K msg/s (80%, gap = pd-ssd) |
| Engine natif (SDK embedded) | N/A | 3.8M evt/s |
| Agents IA intégrés | Non (besoin Kafka Streams + service externe) | Oui (LLM, rules, MCP tools) |
| Graphes d'agents | Non | Oui (DSL fluent, routing conditionnel) |
| Elastic partitioning | Non (partitions fixes) | Oui (16384 slots, rebalance auto) |
| Géo-réplication | MirrorMaker (outil séparé) | Intégrée (active-passive + active-active) |
| io_uring | Non (mmap + page cache) | Oui (SQPOLL, DMA buffers) |
| Déploiement | ZooKeeper 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 StreamFlowCe qui fonctionne sans changement
| Feature Kafka | Supporté | Notes |
|---|---|---|
| KafkaProducer (Java, Python, Go, Node.js) | ✅ | Toutes les libs Kafka |
| KafkaConsumer with consumer groups | ✅ | Rebalancing, offset commit |
| acks=0, acks=1, acks=all | ✅ | Mêmes sémantiques |
| batch.size, linger.ms | ✅ | Mêmes configs |
| compression.type (snappy, lz4, zstd) | ✅ | Décompression transparente |
| Topic auto-creation | ✅ | Configurable |
| kafka-console-producer.sh | ✅ | Outils CLI standard |
| kafka-console-consumer.sh | ✅ | |
| kafka-producer-perf-test.sh | ✅ | |
| kafka-consumer-perf-test.sh | ✅ | |
| Idempotent producer | ✅ | InitProducerId supporté |
| SASL PLAIN authentication | ✅ |
Ce qui n'est pas (encore) supporté
| Feature Kafka | Status | Alternative StreamFlow |
|---|---|---|
| Kafka Transactions | ❌ | acks=all + Raft replication |
| Kafka Streams | ❌ | Agent graphs (plus puissant) |
| Kafka Connect | ❌ | Source/Sink Connector SDK (19 built-in) |
| Schema Registry | ❌ | Headers + validation SDK |
| Exactly-once semantics (EOS) | Partiel | Idempotent producer + at-least-once consumer |
| ACL via kafka-acls.sh | ❌ | REST API + TopicAclPolicy |
| Log compaction | ✅ | LogCompactor 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.
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 KafkaPhase 2 : Dual-write (1-2 semaines)
Vos producers écrivent dans les deux clusters. Les consumers lisent de StreamFlow.
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 :
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 ordersPhase 3 : Basculement (1 jour)
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_publishedPhase 4 : Décommissionner Kafka (après 1 mois)
1# Once StreamFlow is stable in production
2docker compose -f kafka-cluster.yml down
3# or
4kubectl delete statefulset kafka4. Migration par langage
Java (KafkaProducer/KafkaConsumer)
Avant (Kafka) :
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) :
1StreamFlowClient client = StreamFlowClient.remote("streamflow-cluster:9092");
2StreamFlowProducer producer = client.producer();
3producer.send("orders", "key-1", payload);Ou sans changer le code du tout :
1// ONLY CHANGE: the cluster address
2props.put("bootstrap.servers", "streamflow-cluster:9092");
3// The rest of the KafkaProducer code works as-isPython (confluent-kafka)
Avant :
1producer = Producer({'bootstrap.servers': 'kafka-cluster:9092'})Après :
1producer = Producer({'bootstrap.servers': 'streamflow-cluster:9092'})
2# That's it. Same code, same lib.Node.js (KafkaJS)
Avant :
1const kafka = new Kafka({ brokers: ['kafka-cluster:9092'] });Après :
1const kafka = new Kafka({ brokers: ['streamflow-cluster:9092'] });Go (kafka-go)
Avant :
1writer := &kafka.Writer{Addr: kafka.TCP("kafka-cluster:9092"), Topic: "orders"}Après :
1writer := &kafka.Writer{Addr: kafka.TCP("streamflow-cluster:9092"), Topic: "orders"}5. Migration Kafka Connect → StreamFlow Connectors
Kafka Connect (avant)
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)
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=1000Mapping des connecteurs
| Kafka Connect | StreamFlow | Config type |
|---|---|---|
| JdbcSourceConnector | JdbcSourceConnector | jdbc |
| JdbcSinkConnector | JdbcSinkConnector | jdbc |
| FileStreamSourceConnector | HttpSourceConnector | http |
| ElasticsearchSinkConnector | ElasticsearchSinkConnector | elasticsearch |
| MongoSinkConnector | MongoDbSinkConnector | mongodb |
| S3SinkConnector | ❌ (not yet) | — |
| DebeziumSourceConnector | CdcSourceConnector | cdc |
6. Kafka Streams → Agent Graphs
Kafka Streams (avant)
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)
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);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.