The full upstream README, mirrored here for reference. Install config, tool schemas, adoption signals, and an original overview live on the Mq Bridge App listing page.
mq-bridge is an asynchronous message library for Rust. It connects message brokers, databases, files, HTTP/WebSocket endpoints, and in-memory channels behind one small set of traits.
It is not only a forwarder. A route can transform, filter, fan out, retry, rate-limit, deduplicate, or turn a request into a response before the message reaches the next system. The core is built on Tokio and keeps the transport details at the edge, so application code can mostly work with CanonicalMessages and handlers.
If you need to move data or events reliably between systems and you write code (Rust, Python, or Node), mq-bridge is a strong default. It is a library you embed, not a daemon or control plane you operate.
Prefer not to write code? mq-bridge-app runs the exact same engine as a standalone, zero-code ETL service configured entirely by YAML or environment variables — move data from A to B without writing a line. It ships a Postman-style UI to build, send, and inspect messages against a route, and can import Postman collections and AsyncAPI documents to scaffold routes and endpoints for you.
dir_spool (a crash-safe directory FIFO queue), and in-memory channels — all behind the same receive_batch / send_batch shape.pgoutput) and MongoDB (change streams) as flat rows with an operation marker.TlsConfig block (CA bundle, client cert/key for mTLS, insecure-skip) is reused across transports.Throughput & footprint. In our own benchmarks, the same engine — driven the zero-code way through
mq-bridge-app— on a CSV→JSONL file conversion (1,000,000 mixed-type rows, ~116 MiB) sustained 2,824,858 rows/s at ~28 MiB (mq-bridge 0.4.12), about ~145x faster and ~16x leaner in memory than Meltano (tap-csv→target-jsonl, ~19,500 rows/s / ~444 MiB, measured in an earlier session on the same machine). On the same file in the same session it also ran ~1.4x faster than DuckDB using all of this machine's cores (2,036,659 rows/s) while holding ~17x less memory — DuckDB being a throughput ceiling for the conversion itself, not an ETL tool. Methodology and reporting rules are inbenches/ETL_BENCHMARKS.md; the measured numbers and their baselines are in the ETL benchmark harness.
Kafka → file. In a 1,000,000-row relay with the default file format and no transform, the engine was ~65% faster than Sea Streamer comparing both on the mimalloc allocator (~80% against its default-allocator build). The
mq-bridge-appbenchmark contains the reproducible helper and native file-format caveats.
mq-bridge is a Rust library, but the same engine ships as native bindings for Python and Node.js. The Tokio runtime, broker I/O, routing, and batching all stay in Rust; the binding is a thin layer for handlers and configuration.
| Language | Package | Install |
|---|---|---|
| Rust | mq-bridge | cargo add mq-bridge |
| Python | mq-bridge-py (PyPI) | pip install mq-bridge-py |
| Node.js | mq-bridge (npm) | npm install mq-bridge |
The constructor names are kept aligned across languages, so a config loader reads the same in either binding (Python uses snake_case, Node uses camelCase):
Route.from_file / Route.fromFile — load a route from a YAML/JSON fileRoute.from_str / Route.fromStr — load from an in-memory YAML/JSON stringRoute.from_config / Route.fromConfig — load from a dict / JS objectPublisher.* constructors build a publisher endpointThe name argument is optional in both: pass it to select one entry from a routes:/publishers: document, or omit it to treat the config as a single bare route/endpoint body.
The Python binding also holds up well under load on the third-party http-arena.com requests-per-second HTTP benchmark (live leaderboard — rankings shift over time). See the Python analysis notes for the local HTTP comparison harness.
See ARCHITECTURE.md for a detailed overview of the internal design, extensibility, and usage patterns.
Usage Types:
publish / publish_batch and receive / receive_batch directly on endpoints. This mode requires manual commit, batch sequencing, and concurrency handling.For implementation details and quick start examples for each usage type, see the Architecture Guide.
dir_spool (a crash-safe directory FIFO queue), cloud object storage (S3 / GCS / Azure), AWS (SQS/SNS), IBM MQ, and in-memory channels.postgres_cdc) and MongoDB change streams, surfaced as flat rows with an *.operation marker (insert/update/delete/truncate).
Note: IBM MQ is included in the
fullfeature set via theibm-mqfeature, which loads the IBM MQ client library at runtime via dlopen — no IBM SDK is needed to build. The IBM MQ redistributable client only has to be present at runtime, and only if you actually use an IBM MQ endpoint (it is loaded lazily on first connect). If the client is missing, the affected route fails fast with a non-retryable error instead of reconnecting forever. The loader finds the client via the platform default name,MQ_INSTALLATION_PATH(e.g./opt/mqm), or an explicitMQB_IBM_MQ_LIBpath. To link the client statically at build time instead, use theibm-mq-staticfeature (requires the IBM MQ SDK). See the mqi crate for details.TLS: IBM MQ has its own TLS config shape (
IbmTlsConfig), since the native client doesn't consume PEM files. The field names mirror the genericTlsConfigfor config parity, but carry MQ-native semantics:tls.cert_file(aliaskey_repository) is a CMS key repository path (e.g./path/to/tlsfortls.kdb), not a PEM file. The repository can either be passwordless (backed by a.sthstash file next to the.kdb) or password-protected viatls.cert_password(aliaskey_repository_password; requires an IBM MQ client/server at 9.3.0.0 or later, which is the capability levelibm-mq/ibm-mq-staticbuild against; 9.2 is EOL).
send_batch / receive_batch shape. Routes default to batches of up to 512 messages, configurable with batch_size.sled.The project has one main bias: move data reliably without forcing the rest of the application to care too much about the transport.
That means mq-bridge tries to keep the boring parts boring. Kafka offsets, RabbitMQ nacks, HTTP responses, MongoDB polling, WebSocket frames, and file rows are all different in real life, but route code should still be able to receive a batch, process it, publish it, and commit it.
Batching is a big part of that design. Every endpoint is optimized around batch-shaped APIs, even when the backend itself only has a single-message primitive. Routes default to opportunistic batches of up to 512 messages: the consumer waits for the first message, then takes whatever else is already available without adding idle latency. Lower batch_size when tighter failure isolation or a smaller per-batch memory footprint matters.
The error handling follows the same idea. Batch publishing can report partial success, retryable failures, and non-retryable failures. Route commits are sequenced for cumulative-ack brokers so later batches cannot acknowledge earlier unresolved messages; transports with independent acknowledgements can commit batches concurrently. In other words: batching is not just a performance trick bolted onto the side; ack/nack behavior and retry/DLQ handling were built to work with it.
What it does not try to be: a domain framework, an actor runtime, or a full stream processor. You can build CQRS-ish flows with it, but the library cares more about transport, routing, and delivery behavior than about prescribing your domain model.
mq-bridge overlaps with config-driven ETL / pipeline runners, but comes in two shapes: a library you embed in your own service, or mq-bridge-app — the same engine as a zero-code, config-driven ETL service. You get the developer-first path and the no-code path from one codebase.
mq-bridge library when you write the code that moves the data — in Rust, Python, or Node — and want batching, retries/DLQ/dedup, request-reply, CDC, and TLS behind one small API you run in-process. No separate scheduler, no daemon to operate.mq-bridge-app for zero-code ETL when you'd rather move data A→B purely by config: define routes and endpoints in YAML/env, then use its Postman-style UI to send test messages and watch them flow — and import Postman collections or AsyncAPI specs through the UI to generate the config. Same reliability engine, no build step.Either way it stays transport-first on purpose: it moves and routes data reliably rather than prescribing a domain model or a platform.
mq-bridge is young (created in 2025) but its reliability behavior is exercised by an automated integration + performance suite across every supported endpoint:
benches/ETL_BENCHMARKS.md.Known rough edges (be deliberate here):
TlsConfig shape is reused across transports and exercised by automated tests; per-backend certificate matrices (self-signed / CA-signed / mTLS / expired-cert rejection) are being expanded — treat non-trivial TLS setups as needing your own verification for now.CanonicalMessages and swap the underlying transport later.mq-bridge handles the bus, not the entity.mq-bridge intentionally exposes a common subset: publish/consume, pub/sub where possible, request-reply where possible, batching, middleware, and ack/nack handling. If your application depends on highly specific broker features, using that broker's native client directly may be better.input to one output.CommandHandler) or subscribe them (EventHandler).REFERENCE.md lists every middleware and every structural endpoint (
ref,fanout,switch,request,response,reader,static,stream_buffer,null,custom) with fields, defaults and working examples.To add an endpoint of your own, see EXTENDING.md — or PLUGINS.md to ship a Rust endpoint or middleware as a native library that Rust, Python and Node.js hosts all load at runtime, without compiling it into mq-bridge.
mq-bridge endpoints generally default to a Consumer pattern (Queue), where messages are persisted and distributed among workers. To achieve Subscriber (Pub/Sub) behavior, specific configuration is required.
The table below summarizes the capabilities and configuration for each backend:
| Backend | Subscriber Config (Pub/Sub) | Request-Reply | Nack Support |
|---|---|---|---|
| AMQP | Set subscribe_mode: true | Emulated (Property) | Yes (Basic.nack) |
| AWS | N/A (Use SNS) | No | Yes (Visibility Timeout) |
| Directory spool | Set drain_on_read: false (non-destructive fan-out) | No | Yes (chunk left on disk, redelivered) |
| File | Set mode: subscribe | No | Simulated (In-Memory) |
| gRPC | N/A | No | Yes for the built-in Bridge protocol (NACK replays); No for dynamic services unless their API defines an acknowledgement contract |
| HTTP | N/A | Native (Implicit) | Yes (HTTP 500) |
| IBM MQ | Set topic | No | Yes (Tx Rollback) |
| Kafka | Omit group_id | Emulated (Header) | Eventual (Skip Offset) |
| Memory (in-process) | Set subscribe_mode: true | Emulated (Metadata) | Yes (Re-queue), by default disabled |
Memory (IPC: ipc://, unix://, pipe://) | Not supported | Not supported | Yes (Re-queue), by default enabled, consumer-local |
| MongoDB | Set consume: subscriber | Emulated (Metadata) | Yes (Unlock) |
| MQTT | Set clean_session: true | Emulated (Property) | Eventual (Skip Ack) |
| NATS | Set subscriber_mode: true | Native (Inbox) | Yes (JetStream Nak) |
| Postgres CDC | N/A (streams committed changes) | No | Yes (confirmed LSN not advanced) |
| Redis Streams | Set subscriber_mode: true | No | Eventual (PEL, un-acked) |
| Sled | Set delete_after_read: false | No | Yes (Tx Rollback) |
| SQLx | Not supported | No | Eventual (Skip Delete) |
| WebSocket | N/A | No | No |
| ZeroMQ | Set socket_type: "sub" | Native (REQ/REP) | No |
reply_to metadata field) carrying a correlation_id metadata field.Databases have no native pub/sub, so mq-bridge reads them as a source in one of two ways:
postgres_cdc endpoint streams a logical-replication slot (pgoutput). It emits flat JSON rows tagged with postgres.operation metadata (insert/update/delete/truncate). Acking a batch confirms the LSN back to the server (standby_status_update); a nack (or an interrupted run) does not advance the confirmed LSN, so replication resumes from the last acknowledged position on reconnect — at-least-once, verified by a restart-safety integration test. Enable with the postgres-cdc feature; requires wal_level = logical and a publication. See below.consume: capture_new to watch an existing collection for changes from now on, or consume: capture_all to read the existing documents first and then keep capturing changes (no gap, at-least-once). It emits insert/update/replace/delete events tagged with mongodb.operation metadata and checkpoints progress under cursor_id. Both use the change stream (full post-image via updateLookup) and therefore need a replica set — a single-node one is enough. Without one they refuse to start; use consume: snapshot for a one-shot non-destructive read, or consume: consumer for a work queue.cursor_column (WHERE col > $last ORDER BY col ASC), persisting the last read value under cursor_id. Captures appends only — updates and deletes are not observed. Available on SQLx (PostgreSQL / MySQL / MariaDB / SQLite) and ClickHouse; while idle the poll interval backs off exponentially between polling_interval_ms and max_polling_interval_ms. SQLite and ClickHouse are polling-only — they have no server-side change log.SQLx read mode is chosen by config, not by driver. With no
cursor_column, an SQLx source (?table=/table:) is a competing-consumers work queue: it atomically claims rows via alocked_untillease and deletes them on ack, so several readers drain the same table as competing consumers — each row goes to one reader at a time, at-least-once (a lease that expires before ack, e.g. after a crash, redelivers the row), not exactly once. This requires the table to carry the queue schema — anid, apayload, and alocked_untilcolumn — on every driver (PostgreSQL, MySQL/MariaDB, and SQLite behave identically here); it is not a plain full-table read. To read an existing or arbitrary table non-destructively (the ETL path), setcursor_columnto switch to the cursor polling above. A source table missinglocked_untilfails fast with a permanent error rather than reconnecting forever.
These are not interchangeable. The default capture_all is a reader, like the other capture_*
modes: they fan out, so every reader sees every document, and nothing is claimed or removed. Opting
into consumer instead gives a competing-consumers work queue: its lock-and-delete protocol lets
several readers drain the same collection with each document going to exactly one of them — the right
thing for dispatching jobs or commands, and destructive by design.
Pick by semantics first. Where both would do — a one-shot bulk read or ETL pass with a single
reader — use capture_all, the default: consumer's four round trips per batch (find ids →
claim → re-fetch → delete on ack) buy exclusivity you aren't using. Set
consume: consumer explicitly when you want the destructive work-queue semantics.
consume | mechanism | modifies source | ends on drain | use for |
|---|---|---|---|---|
capture_all (default) | snapshot, then change stream | no | with drain options | bulk read / ETL |
consumer | claim → lock → re-fetch → delete | yes | yes | work queues, competing readers |
snapshot | one _id-ordered pass over the collection | no | always | reads without a replica set |
capture_new | change stream, new changes only | no | no | ongoing CDC |
500k documents, batch_size: 1024, release build, standalone MongoDB: consumer → null
23,667 rows/s (postgres → null reference: 115,924). The previously quoted capture_all
figure was measured on that same standalone server, i.e. through the _id fallback that has since
been removed as unsound — it is not comparable and needs re-measuring on a replica set.
concurrency does not speed up a MongoDB source — batches are fetched serially; it only widens
the downstream side.consumer only reads collections written by the mq-bridge MongoDB publisher (UUID _id /
wrapped envelope); other documents are skipped or rejected at startup. snapshot and the
capture_* modes read arbitrary collections.capture_* need a replica set — a single-node one is enough — and both error without it. With
exit_on_empty / --drain, a replica-set reader finishes after the snapshot and any
immediately following changes, once the drain idle timeout expires.snapshot is one-shot by design: it delivers what exists when the run starts and then ends the
route. It rejects cursor_id, because resuming above a stored _id would skip anything a
concurrent writer commits below that mark — see the note below.capture_all's snapshot also emits the internal <collection>:sequencer document.Why
snapshotdoes not resume. A_idcursor can only match documents above its high-water mark, and_idis assigned client-side before the insert, so it does not follow commit order. Any writer that commits below the mark — concurrent producers, a route withconcurrency > 1, a retried write — is skipped permanently and silently. Incremental reads therefore need commit order, which on MongoDB means the oplog, which means a replica set. There is no sound_id-based substitute.
With a durable source and durable checkpoint configuration, mq-bridge is at-least-once across
crashes: a replay can redeliver, while in-process Memory endpoint state is not crash-durable. Pair a
stable replay identity with an idempotent sink operation and that covered sink effect is
effectively exactly-once, however often the message is delivered. Which sink absorbs a duplicate, which source gives you a stable key to
deduplicate on, and what a handler in the route changes about all of this, is covered in full by:
docs/DELIVERY.md — delivery guarantees. Per-source identity and per-sink idempotency tables, the
deduplicationmiddleware, Mongoid_fieldandreport_outcome, SQLON CONFLICT, ClickHouseReplacingMergeTree, file/object-storename_by: source_position, what handlers change, and whatmq-bridgedeliberately does not provide.
The short form: point the sink at a deterministic business key (id_field for MongoDB, an
ON CONFLICT column for SQL) and let its unique constraint be the authority — it is already shared
across every writer, so no extra state store is needed. Add the deduplication middleware on the
input only when the sink has no constraint to lean on, or to keep duplicates away from a handler.
The object_store endpoint (alias s3) reads and writes cloud object stores — Amazon S3, Google Cloud Storage, Azure Blob, Cloudflare R2, and anything else the object_store crate speaks — behind the same receive_batch / send_batch API. Enable it with the object-store feature. Credentials and backend options are read from the process environment (AWS_ACCESS_KEY_ID, AWS_ENDPOINT, AWS_REGION, GOOGLE_SERVICE_ACCOUNT, AZURE_STORAGE_ACCOUNT, ...); the URL scheme picks the backend (s3://, gs://, az://).
normal JSONL, json, text, raw) and written as whole immutable objects — write-once, nothing appended or mutated. Under name_by: write_time each flushed batch becomes one object at <prefix>/[YYYY/MM/DD/]<uuidv7>.<ext>; the uuidv7 name already sorts by write time, and the optional date_partition prefix (on by default, derived from that same id's timestamp) is a readability / lifecycle-rule convenience. Under name_by: source_position a batch becomes one object per contiguous run of source positions it covers — usually one, more when the batch spans partitions or skips offsets a replay already wrote — each named for the range it covers and written flat under the prefix, where date_partition does not apply. The default name_by: auto picks source_position whenever the route's input carries a replay position — Kafka, Postgres CDC, a SQL cursor, a MongoDB change stream or a file — so key order equals source order at any concurrency without configuring anything (see Files & object storage — name_by).cursor_id and an external checkpoint_store (file://, s3://, postgres://, mongodb://) so a restart resumes without re-emitting. Objects are never deleted or rewritten — resume is non-destructive and at-least-once at object granularity (a nacked batch is redelivered; the cursor only advances once an object is fully acked). csv is supported on the source only.Point
checkpoint_storeat a different bucket or prefix than the source reads; a cursor object written under the source prefix would be listed and re-read as data (the source rejects an overlapping object-store checkpoint location).
The sink's
formatdecides the encoding, not where the message came from:normal/json/textalways write the{message_id, payload, metadata}wrapper, so the message id survives the round trip. Useformat: rawto write payloads verbatim (bare documents, no wrapper). Applies to bothfileandobject_store.A source on those formats expects that same wrapper,
message_idincluded. A line that is valid JSON but not the wrapper — a hand-written fixture with onlypayloadandmetadata, say — is not decomposed: the whole line becomes the payload and the line's ownmetadatais discarded. The reader logs a warning naming the wrapper when this happens. For plain JSON lines, useformat: raw.Payload encoding in the wrapper. A UTF-8 payload is written as a plain JSON string under
payload. A binary one (compressed, encrypted, Protobuf, …) is base64-encoded under a separatepayload_base64field; the two are mutually exclusive, as in the CloudEvents JSON format.Sources still read the older byte-array form (
"payload":[123,34,…]), so existing files keep working. Only the reverse is a break: a binary record written by this version is not readable by an older mq-bridge. Text records are compatible in both directions.
The response output endpoint sends a reply back to the original requester. This is useful for synchronous request-reply flows, for example HTTP-to-NATS-to-HTTP. Use response: {} as the output endpoint configuration.
response will be dropped.correlation_id) may break the response chain.The mq-bridge-app workspace packages run mq-bridge as a
standalone CLI, desktop app, or Docker container configured via YAML or
environment variables. Plain cargo build still builds only this library;
use cargo build --workspace to build every workspace package.
mq-bridge-app can be used to create and test route and endpoint configurations through a UI. The generated JSON/YAML can then be copied into an application and loaded by the Rust library or the Python bindings.
This does not replace application code or handlers, but it is useful when you want a known-good connection and route shape before pasting the configuration into code. For Python and Node projects, routes and publishers are commonly loaded from JSON/YAML with Route.from_file/fromFile (a file path) and Route.from_str/fromStr (a JSON/YAML string), while Route.from_config/fromConfig takes an already-parsed dict/JS object. The matching Publisher.* constructors work the same way.
For business logic, mq-bridge provides a handler layer separate from transport-level middleware. This is where message-specific code usually belongs.
CommandHandler: A handler for 1-to-1 or 1-to-0 message transformations. It takes a message and can optionally return a new message to be passed down the publisher chain.EventHandler: A terminal handler that reads new messages without removing them for other event handlers.You can chain these handlers with endpoint publishers.
For more structured, type-safe message handling, mq-bridge provides TypeHandler. It deserializes messages into a specific Rust type before passing them to a handler function, so handlers do not need to repeat the same parsing code.
Message selection is based on the kind metadata field in the CanonicalMessage.
You can define and run routes directly in Rust code.
mq-bridge supports request-response patterns for interactive services such as web APIs. A client can send a request and wait for the matching response, while the bridge keeps the correlation details away from the handler.
The response output is the most direct option and the safest one under concurrency.
response Output Endpoint (Recommended)For request-response routes, use the dedicated response endpoint in the route's output.
How it works:
http) receives a request.handler to process the request and generate a response payload.output.response: {}, the bridge sends the message back to the original input source, which then sends it as the reply (e.g., as an HTTP response).The response stays in the same execution context as the request, so concurrent requests do not need to share a reply queue and race on correlation IDs.
For example, a service can write a request document to MongoDB and wait for a reply. The bridge reads the document, runs the handler, and writes the result back to the reply collection.
YAML Configuration (mq-bridge.yaml):
Programmatic Handler Attachment (in Rust): You would then load this configuration and attach a handler to the route's output endpoint in your Rust code.
mq-bridge can be used for CQRS-style flows. With routes and typed handlers, it can act as a command bus and an event bus without becoming a domain framework.
All routes and endpoints can be defined via a configuration file (for example mq-bridge.yaml), JSON, or environment variables. For a complete reference of options, middleware, and examples, see the Configuration Guide.
Important route-level knobs:
batch_size: maximum messages per route iteration. Defaults to 512; lower it for tighter failure isolation or large payloads.concurrency: number of route workers. Defaults to 1; useful for high-latency handlers or endpoints.commit_concurrency_limit: maximum in-flight commit operations, whether they are queued through ordered commit sequencing or run concurrently for independent-ack transports. Defaults to 4096.Middleware can be attached to inputs or outputs. The most commonly used ones are retry, dlq, deduplication, limiter, cookie_jar, and transform. Retry/DLQ are especially useful with batching because partial failures can be retried or sent to a DLQ without treating the entire batch as equally broken.
The transform middleware reshapes JSON payloads declaratively, so field renaming and type
fixing do not need a custom handler. It runs two optional stages over a single parse:
mapping (rename, move, nest) and then schema (coerce, apply defaults, validate).
With the CSV source above producing {"first_name":"John","last_name":"Smith","user_id":"42"},
the mapping yields {"firstName":"John","lastName":"Smith","id":"42","address":{"city":"unknown"}}
($.city is absent, so address.city falls back to its "unknown" default) and the schema
then coerces id to the integer 42.
Failures are always non-retryable and name the offending field, e.g.
transform failed at $.items[1].qty [coercion]: cannot coerce string "oops" to integer.
On an output endpoint the message is failed so a following dlq captures it; on an
input endpoint it is dropped from the batch and acknowledged, which keeps invalid data
out of the route.
See REFERENCE.md for the full option list, the supported schema subset, and the exact set of allowed type coercions.
The project includes integration and performance tests. Most backend tests require Docker.
To run the performance benchmarks for all supported backends:
To run the criterion benchmarks:
The times are not stable yet, it is therefore recommended to perform the integration performance test if you want to measure throughput.
Contributions are welcome. See CONTRIBUTING.md for setup notes, code style, and pull request guidelines.
This library has been written with a lot of AI assistance.
The core started as my own code, and many endpoints and docs were expanded with help from Gemini, CodeRabbit, Claude, and Codex. The useful part was speed: once the endpoint traits were stable, adding more transports became much easier. The dangerous part is the usual one: generated code can look plausible while still missing important details. I am aware that in year 2026, AI is still not generating perfect code and sometimes breaks simple stuff or forgets important lines during refactorings that later result in severe bugs.
For that reason I reviewed each commit manually to prevent hard-to-fix architectural issuess and cleaned up, and refactored the generated output.
I do trust the current code as much as if it would be written by myself.
Due to the large feature set, there may still be unfixed issues. The current focus is testing and documentation.
Some parts of the code are more verbose than I would write by hand, but I kept the readable parts when they worked well. I am not a native English speaker, so AI assistance is also useful for documentation. The important part is that the code is reviewed and tested, not that every sentence or helper function looks hand-typed from the first draft.
mq-bridge is licensed under the MIT License.
Binary distributions carry their applicable third-party terms in
THIRD_PARTY_LICENSES.txt. Their current maintenance process is documented in
docs/THIRD_PARTY_LICENSES.md.