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.

Diagram of the pipeline data flow inside a single JVM process: a Pulsar source feeds the PulsarSource connector, then the STRING_TO_WIREEVENT transform, then a Kafka sink.
Pipeline data flow — in-process, single JVM

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

ParameterValue
Pipeline IDsk-pulsar-source-kafka
Parallelism4 worker threads
Batch size128 records
SourcePULSAR — persistent://public/default/streamkernel-bench-in
Subscriptionstreamkernel-pulsar-kafka · Exclusive · EARLIEST
TransformSTRING_TO_WIREEVENT — payload ID as routing key
SinkKAFKA — streamkernel-pulsar-out · 12 partitions
Kafka producerlz4 compression · 128KB batch · 256MB buffer · acks=1
ProvenanceEnabled — per-event lineage stamped on all output records
MetricsPrometheus on :8080

03 Test Methodology

3.1 Environment

ComponentDetails
Host machineSEKISOFT — Windows 10 (WSL2 Docker) — consumer laptop
JVM heap6GB fixed (–Xms6g –Xmx6g) · G1GC · MaxGCPauseMillis=50
StreamKernelv0.2.0 — streamkernel-app-0.2.0-all.jar
Apache PulsarStandalone Docker container — localhost:6650
Apache KafkaConfluent Platform 7.6.1 — KRaft mode — localhost:9092
Kafka partitions12
Run window20: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:

Bar chart of processed and output records per second in four consecutive windows during the burst phase.
Speedometer windows — records/sec during burst phase (5-second intervals)

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:

Bar chart of the number of records in each of the twelve Kafka output partitions, p0 to p11.
Kafka partition distribution — streamkernel-pulsar-out · 12 partitions · 253,235 total records

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:

CapabilityStatusNotes
Transport portability✓ COMPLETESame JAR, SPI-driven, config-only swap
Record integrity (functional)✓ COMPLETEZero loss across five independent measurement points
Per-event provenance / lineage✓ COMPLETEStamped on all output records via WireEvent envelope
Prometheus observability✓ COMPLETEFull counter/gauge coverage, consistent label set
Delivery guarantee documentationCONFIG ONLYRuntime supports at-least-once; acknowledge.on.fetch=false path available and validated in other profiles
Dead-letter queueCONFIG ONLYDEVNULL excluded from benchmark to isolate throughput; Kafka DLQ sink validated in production profiles
Cross-transport atomicityCONFIG ONLYAck-after-confirm path available in runtime; excluded from this benchmark intentionally
Security (mTLS + OPA)CONFIG ONLYPERMIT_ALL excluded from benchmark; mTLS + OPA_SIDECAR profiles validated independently
Idempotent Kafka producerCONFIG ONLYacks=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

ParameterBenchmarkProduction recommendation
JVM heap6GB4–8GB for most bridging workloads at this message size
Worker parallelism4 threads8–16 for server-class hardware
Kafka producer buffer256MBIncrease for higher sustained throughput targets
Expected throughput (2× parallelism, dedicated brokers)15,647 r/s peakConservatively 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

Commercial path

Want to turn this paper into a concrete evaluation?

Bring the action you have in mind and we will map it to what the runtime does today.