Research library / Whitepaper
Pulsar-to-Kafka Transport Bridging
Cross-transport pipeline portability via a single JVM runtime — no code changes, no managed connectors.
- Pipeline Profile
- sk-pulsar-source-kafka
- Run ID
- run-pulsar-source-01
- JAR
- streamkernel-app-0.2.0-all.jar
- Patent
- US Prov. 64/057,035
00 Executive Summary
Enterprise messaging infrastructure is rarely homogeneous. Organizations routinely operate Apache Pulsar and Apache Kafka in parallel — across divisions, cloud regions, migration phases, or compliance boundaries — and require a reliable, operationally simple bridge between the two.
StreamKernel addresses this with a single-JAR JVM runtime that consumes from Pulsar and produces to Kafka entirely in-process, with no external orchestration, no intermediate storage layer, and no code changes required when swapping transports. This document presents the architecture, configuration, and empirical performance results of a validated Pulsar-to-Kafka pipeline run on a development laptop.
- Peak throughput
- 15,647 records / second
- Total records
- 253,235 zero loss · zero drops
- Burst duration
- ~20s to drain 253K backlog
- GC overhead
- 0.1% 604ms / 623s · no major GC
01 The Enterprise Problem
Several documented customer scenarios drive demand for Pulsar-to-Kafka bridging. A Technical Support Engineer at a major platform vendor has confirmed this as an active customer requirement.
Messaging infrastructure consolidation
Organizations that adopted Pulsar for its multi-tenancy or geo-replication properties and are now standardizing on Kafka-based data platforms need a drain path. Rewriting producers is costly and risky; a bridge runtime is the lowest-friction option.
Cross-organizational event handoff
A business unit or partner system running Pulsar needs to feed events into a central Kafka platform for analytics, ML pipelines, or downstream consumers. This boundary is often a compliance or ownership boundary as well, making a standalone bridge preferable to shared infrastructure.
Air-gapped and regulated deployments
Defense, government, and financial environments frequently segment networks by classification level or regulatory domain. A single-JAR runtime that operates entirely in-process — with no Kubernetes dependency, no managed connectors, and no external state — fits these environments in ways that cloud-native integration platforms do not.
Incremental migration
Organizations migrating from Pulsar to Kafka cannot cut over all producers simultaneously. StreamKernel allows legacy Pulsar topics to continue receiving data while the migration proceeds, draining them into Kafka without altering producer behavior.
02 Architecture
2.1 Design philosophy
StreamKernel is built around a strict SPI (Service Provider Interface) abstraction that decouples source transport, transform chain, and sink transport at the plugin boundary. Swapping Pulsar for Kafka as source — or Kafka for MongoDB as sink — requires only a properties file change. No recompilation. No new binary.
2.2 Plugin catalog
StreamKernel's SPI discovery loaded the full plugin catalog at startup, confirming the breadth of available transports within this single binary:
# Plugin catalog — discovered at boot, same fat JAR
Sources: [KAFKA, PULSAR, REST, SALESFORCE, SYNTHETIC]
Sinks: [DELTA, DEVNULL, KAFKA, MONGO_INSERT, MONGO_VECTOR, SNOWFLAKE_SNOWPIPE_STREAMING]
Transformers: [DETERMINISTIC_ENRICHMENT, EMBEDDING_TO_ENRICHED_TICKET, HTTP_EMBEDDING, NOOP, STRING_TO_WIREEVENT]
Security: [OPA_SIDECAR, PERMIT_ALL]
2.3 Pipeline configuration
| Parameter | Value |
|---|---|
| Pipeline ID | sk-pulsar-source-kafka |
| Parallelism | 4 worker threads |
| Batch size | 128 records |
| Source | PULSAR — persistent://public/default/streamkernel-bench-in |
| Subscription | streamkernel-pulsar-kafka · Exclusive · EARLIEST |
| Transform | STRING_TO_WIREEVENT — payload ID as routing key |
| Sink | KAFKA — streamkernel-pulsar-out · 12 partitions |
| Kafka producer | lz4 compression · 128KB batch · 256MB buffer · acks=1 |
| Provenance | Enabled — per-event lineage stamped on all output records |
| Metrics | Prometheus on :8080 |
03 Test Methodology
3.1 Environment
| Component | Details |
|---|---|
| Host machine | SEKISOFT — Windows 10 (WSL2 Docker) — consumer laptop |
| JVM heap | 6GB fixed (–Xms6g –Xmx6g) · G1GC · MaxGCPauseMillis=50 |
| StreamKernel | v0.2.0 — streamkernel-app-0.2.0-all.jar |
| Apache Pulsar | Standalone Docker container — localhost:6650 |
| Apache Kafka | Confluent Platform 7.6.1 — KRaft mode — localhost:9092 |
| Kafka partitions | 12 |
| Run window | 20:35:10 UTC → 20:45:33 UTC (10.38 min) |
Note on test conditions: Both brokers ran in Docker containers on the same laptop as the StreamKernel JVM, sharing CPU and memory. No network isolation, no dedicated hardware. This is a conservative baseline — production server hardware with dedicated brokers will yield substantially higher sustained throughput.
3.2 Test execution
A purpose-built PowerShell pre-run script (demo_before_pulsar_source_kafka.ps1) handled full environment preparation — including Pulsar topic reset with BookKeeper ledger recovery on error, deterministic backlog seeding, and subscription creation at EARLIEST with confirmed backlog verification before pipeline start.
3.3 Correctness baseline
A prior run with a 10,000-message backlog confirmed end-to-end record integrity before scaling up: 10,000 records processed, DROPPED=0, errorTotal=0, Pulsar deliverCounter=10,000, Kafka total records 10,000 — exact match. The 253,235-message run below is the primary performance evidence.
04 Results
4.1 Throughput
The pipeline consumed the full 253,235-message backlog during the burst phase beginning at pipeline start. Four consecutive speedometer windows (5-second intervals) captured throughput ramp-up from cold start to full speed:
Peak throughput of 15,647 records/second at window 2. The entire 253,235-message backlog was drained in approximately 20 seconds on a consumer-grade laptop. The Pulsar consumer stats recorder independently confirmed a 60-second average consume throughput of 4,220 msgs/s — consistent with the burst-then-idle pattern.
- Peak PROC_EPS
- 15,647 window 2 of burst
- Backlog drain time
- ~20s 253,235 records
- Pulsar avg. ack rate
- 4,220 msgs/s (60s avg)
4.2 Record integrity — five-point verification
End-to-end record integrity was verified independently across three measurement layers: Pulsar broker telemetry, StreamKernel Prometheus counters, and Kafka physical partition counts. All values agree exactly.
- Pulsar publish
- 253,235
- Pulsar deliver
- 253,235
- SK processed
- 253,235
- Kafka sink sent OK
- 253,235
- Kafka physical
- 253,235
From Prometheus snapshot (run-pulsar-source-01): streamkernel_pipeline_dropped_total = 0 · streamkernel_pipeline_source_errors_total = 0 · streamkernel_pipeline_dlq_total = 0 · streamkernel_pipeline_auth_errors_total = 0. Pulsar subscription drained to backlog=0, unacked=0.
4.3 Kafka partition distribution
Records were distributed across all 12 output partitions via key-based routing from the STRING_TO_WIREEVENT transform, confirming effective routing and balanced producer behavior:
4.4 Transform performance
The STRING_TO_WIREEVENT transformer processed all 253,235 records with negligible per-record cost:
- Avg. encode time
- ~1µs per record
- Encoded payload
- 119B per WireEvent
- Kafka req latency
- 51.4ms avg (co-located broker)
- Buffer utilization
- <1% of 256MB producer buffer
Cumulative encode time: 253.285 milliseconds across 253,235 records = ~1µs average. The transform step is effectively zero-cost relative to network I/O — consistent with the design goal of in-process, allocation-minimal transformation. Buffer utilization below 1% confirms the Kafka producer was never backpressured during the run.
4.5 JVM & GC behavior
G1GC performed 10 pause events across the full run. No major GC cycles occurred. Old generation promotion was zero bytes.
- G1 Evacuation (×8)
- 536 ms
- Metadata GC (×2)
- 68 ms
- Total pause / runtime
- 0.1%
- Post-GC heap live data
- 32MB of 6,144MB heap (0.52%)
- Old gen promotion
- 0B no long-lived accumulation
- Peak thread count
- 27 12 live at snapshot
05 Enterprise Readiness
The results above demonstrate functional correctness and strong burst throughput. The following assessment maps current capability against enterprise production requirements:
| Capability | Status | Notes |
|---|---|---|
| Transport portability | ✓ COMPLETE | Same JAR, SPI-driven, config-only swap |
| Record integrity (functional) | ✓ COMPLETE | Zero loss across five independent measurement points |
| Per-event provenance / lineage | ✓ COMPLETE | Stamped on all output records via WireEvent envelope |
| Prometheus observability | ✓ COMPLETE | Full counter/gauge coverage, consistent label set |
| Delivery guarantee documentation | CONFIG ONLY | Runtime supports at-least-once; acknowledge.on.fetch=false path available and validated in other profiles |
| Dead-letter queue | CONFIG ONLY | DEVNULL excluded from benchmark to isolate throughput; Kafka DLQ sink validated in production profiles |
| Cross-transport atomicity | CONFIG ONLY | Ack-after-confirm path available in runtime; excluded from this benchmark intentionally |
| Security (mTLS + OPA) | CONFIG ONLY | PERMIT_ALL excluded from benchmark; mTLS + OPA_SIDECAR profiles validated independently |
| Idempotent Kafka producer | CONFIG ONLY | acks=1 excluded from benchmark; acks=all + enable.idempotence validated in hardened profiles |
Benchmark configuration note: The five CONFIG ONLY items above are fully built capabilities within the StreamKernel runtime, validated in other pipeline profiles (mTLS+OPA, hardened Kafka, DLQ). They are deliberately excluded from this benchmark to isolate raw transport throughput. Activating any of them requires only a properties file change — no code, no recompilation, no new binary.
06 Deployment Considerations
Single-JAR deployment model
StreamKernel ships as a self-contained fat JAR. No Kubernetes operator. No sidecar. No external service dependency beyond the source and sink brokers. This makes it suitable for air-gapped environments, edge deployments, and regulated infrastructure where managed connector frameworks are not permitted.
Sizing guidance
| Parameter | Benchmark | Production recommendation |
|---|---|---|
| JVM heap | 6GB | 4–8GB for most bridging workloads at this message size |
| Worker parallelism | 4 threads | 8–16 for server-class hardware |
| Kafka producer buffer | 256MB | Increase for higher sustained throughput targets |
| Expected throughput (2× parallelism, dedicated brokers) | 15,647 r/s peak | Conservatively 30,000–50,000 r/s for this message profile |
Key Prometheus metrics to monitor
# Alert on these in production
streamkernel_pipeline_dropped_total # should always be 0
streamkernel_pipeline_source_errors_total # source connectivity
streamkernel_pipeline_dlq_total # processing failures
streamkernel_kafka_sink_sent_ok_total # lag vs source_pulsar_read_total
streamkernel_kafka_sink_request_latency_avg_ms # broker health proxy
07 Summary
StreamKernel's Pulsar-to-Kafka pipeline profile demonstrates that enterprise-grade cross-transport message bridging can be delivered as a single JVM process with no external orchestration. On a development laptop with co-located Docker brokers:
- Peak throughput
- 15,647 records / sec
- Records processed
- 253,235 zero loss · zero drops
- GC overhead
- 0.1% no major GC
- Transform cost
- ~1µs per record avg
The transport portability proof: The same fat JAR, the same SPI-driven architecture, and the same benchmark harness validated across Kafka, Pulsar, MongoDB, Delta Lake, Snowflake, and mTLS+OPA profiles. Transport is a configuration choice, not a code change.
StreamKernel is source-available software developed by IntuitiveDesigns.
Core runtime: StreamKernel Source Available License (SSAL v1.0) · SPI interfaces: Apache 2.0
Commercial licensing available: Professional · Enterprise · OEM · Managed Service · Government
Contact: steven.lopez@streamkernel.io
US Provisional Patent No. 64/057,035 — Filed May 4, 2026