惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Martin Fowler
Martin Fowler
CTFtime.org: upcoming CTF events
CTFtime.org: upcoming CTF events
A
About on SuperTechFans
Apple Machine Learning Research
Apple Machine Learning Research
The Register - Security
The Register - Security
Vercel News
Vercel News
H
Hackread – Cybersecurity News, Data Breaches, AI and More
人人都是产品经理
人人都是产品经理
MyScale Blog
MyScale Blog
云风的 BLOG
云风的 BLOG
博客园_首页
U
Unit 42
T
Tailwind CSS Blog
G
GRAHAM CLULEY
F
Full Disclosure
V
Vulnerabilities – Threatpost
T
Tenable Blog
月光博客
月光博客
P
Privacy & Cybersecurity Law Blog
P
Privacy International News Feed
K
Kaspersky official blog
Scott Helme
Scott Helme
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
N
News and Events Feed by Topic
T
The Exploit Database - CXSecurity.com
N
News and Events Feed by Topic
有赞技术团队
有赞技术团队
Recent Commits to openclaw:main
Recent Commits to openclaw:main
L
LINUX DO - 最新话题
Recorded Future
Recorded Future
Application and Cybersecurity Blog
Application and Cybersecurity Blog
Help Net Security
Help Net Security
The GitHub Blog
The GitHub Blog
Cisco Talos Blog
Cisco Talos Blog
SecWiki News
SecWiki News
P
Proofpoint News Feed
Security Latest
Security Latest
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
罗磊的独立博客
S
Security Affairs
M
MIT News - Artificial intelligence
L
LINUX DO - 热门话题
美团技术团队
Simon Willison's Weblog
Simon Willison's Weblog
T
Threat Research - Cisco Blogs
Stack Overflow Blog
Stack Overflow Blog
Forbes - Security
Forbes - Security
Hugging Face - Blog
Hugging Face - Blog
博客园 - Franky
V
Visual Studio Blog

DEV Community

Authentication Security Deep Dive: From Brute Force to Salted Hashing (With Java Examples) Why AI Systems Don’t Fail — They Drift Spilling beans for how i learn for exam😁"Reinforcement Learning Cheat Sheet" I Replaced Chrome with Safari for AI Browser Automation. Here's What Broke (and What Finally Worked) How Python Borrows Other People's Work Top 15 Reinforcement Learning Questions That Will Appear in Exams The $40 Architecture: Processing 1 Billion API Requests with 99.99% Uptime Vibe Coding: A Workflow Guide (From Zero to SaaS) Most webhook security guides protect the wrong side. The scary part is delivery. Headless CMS for TanStack Start: Build a Blog with Cosmic EU Age Verification App "Hacked in 2 Minutes" — What Actually Happened Comfy Cloud’s delete function does not actually remove files Running AI Models on GPU Cloud Servers: A Beginner Guide Event-driven media intelligence with AWS Step Functions and Bedrock I scored 500 AI prompts across 8 quality dimensions — here's what broke How to Call Google Gemini API from Next.js (Free Tier, No Backend Needed) The Portal Protocol: Reclaiming Human Connection in the Age of AI How to Fix Your Team's Scattered Knowledge Problem With a Self-Hosted Forum Intro to tc Cloud Functors: A Graph-First Mental Model for the Modern Cloud Designing Multi-Tenant Backends With Both Ownership and Team Access I Built a Neumorphic CSS Library with 77+ Components — Here's What I Learned PostgreSQL Performance Optimization: Why Connection Pooling Is Critical at Scale Cómo construí un SaaS multi-rubro para gestionar expensas en Argentina con FastAPI + Vue 3 🚀 I Built an Ethical Hacking Scanner Tool – Open Source Project I Replaced /usage and /context in Claude Code With a Single Statusline A Pythonic Way to Handle Emails (IMAP/SMTP) with Auto-Discovery and AI-Ready Design I Collected 8.9 Million Polymarket Price Points — Here's What I Found About How Markets Really Move EcoTrack AI — Carbon Footprint Tracker & Dashboard Everyone's Using AI. No One Agrees How. 5 self-hosted ebook managers worth trying in 2026 Building Your First AI Agent with LangChain: From Chatbot to Autonomous Assistant Common SOC 2 Failures (Real World) Stop Vibe-Checking Your AI App: A Practical Guide to Evals How to Use SonarQube and SonarScanner Locally to Level Up Your Code Quality Your Next To-Do App Is Dead — I Replaced Mine with an OpenClaw AI Sign a Nostr event in 60 lines of Python using coincurve — no nostr-sdk, no nbxplorer, no rust toolchain ITGC Audit Explained Like You’re in Big 4 Patch Tuesday abril 2026: Microsoft parcha 163 vulnerabilidades y un zero-day en SharePoint Stop scraping everything: a better way to track competitor price changes Listing on MCPize + the Official MCP Registry while routing payments OUTSIDE the marketplace — how I kept 100% of my x402 revenue Building an AI-Powered Risk Intelligence System Using Serverless Architecture Why We Ripped Function Overloading Out of Our AI Toolchain Testing AI-Generated Code: How to Actually Know If It Works SaaS Churn Is Killing Your Business. Here Is What to Do About It (Without a Support Team) The Speed of AI Is No Longer Linear - And Self-Improving Models Are Why How to Implement RBAC for MCP Tools: A Practical Guide for Engineering Teams From Standard Quote to Persuasive Proposal: AI Automation for Arborists I built a CLI that scaffolds complete multi-tenant SaaS apps Axios CVE-2025–62718: The Silent SSRF Bug That Could Be Hiding in Your Node.js App Right Now The dashboard that ended our friendship Data Pipelines Explained Simply (and How to Build Them with Python) The Hidden Cost of AI Systems Nobody Talks About. undefined vs undeclared, and how typeof behaves io_uring Adventures: Rust Servers That Love Syscalls Why Agentic AI is Killing the Traditional Database The POUR principles of web accessibility for developers and designers Quantum Neural Network 3D — A Deep Dive into Interactive WebGL Visualization How To Install Caveman In Codex On macOS And Windows Automation Pipeline Reliability: Why Your Workflow Breaks When Nobody Is Watching I Built an 'Open World' AI Coding Agent — It Works From ANY Folder From Freelancing to Product: A Tech Service Company's SaaS Transformation China's AI Giants: Adding Tencent Hunyuan & ByteDance Doubao to AI University (74 Providers) On the Vibe Coders and Their Lies clerk: Auto-Summarize Your Claude Code Sessions AI Weekly — 2026/04/10–04/17 | The Model Lockdown Is Here, but the Toolchain Is the Real Battleground AI 週報 — 2026/04/10–2026/04/17 模型封鎖潮來了,但工具鏈才是真戰場 Maybe this is how Open-Source apps are born... 🚀 Fine-Tune LLMs with LoRA and QLoRA: 2026 Guide tRPC v11 + Next.js App Router: End-to-End Type Safety Without the Boilerplate ShadCN UI in 2026: Why I Stopped Installing Component Libraries and Started Owning My Components SaaS Billing in React Server Components: Stripe + Supabase Without a Single `useEffect` Join our DEV Weekend Challenge — $1,000 in Prizes Across TEN winners! Submissions Due April 20 at 6:59 AM UTC. Implementing FSRS Spaced Repetition in Flutter + Supabase — Adding Memory Science to an AI Learning App "I Texted My Localhost From the Train — Claude Code Fixed the Bug Before I Got Home" I Built a Sales Prep AI and It Went Deeper Than Expected Design to Code #2: One JSON, Eleven Outputs Solving the 100M-Row Problem: A Summary Table Pattern for High-Volume Push Notification Logs Flutter Web With Wasm: What Actually Changes For Developers I Built 50 Royalty-Free Soundtracks for My Side Project in a Weekend Using AI Music Generation The Vibe Coding Security Checklist: 7 Things to Check Before You Ship Stop Letting Googlebot Guess Fix Your React App's SEO Right Desconstruindo o Streaming do LinkedIn: Como Criar um Engine de Extração de Vídeo de Alta Performance com HLS e FFmpeg (EDA Part-1) EDA (Exploratory Data Analysis) Explained With Real Life — Why Looking at Your Data Is the Most Important Step in Machine Learning Brand Relationship Management at Scale: Our 4-Touch Outreach System for 200+ Brands Why String.fromEnvironment() Might Return an Empty String in Dart JGuardrails 1.0.0 — Hardening Java LLM Apps Against Jailbreaks, Toxicity, and Prompt Injection Plan and Schedule a Full Week of Threads Content From One Claude Conversation Coding Cat Oran Ep3, Five Tables Changed Everything Updated: BFF Pattern I'm done watching freelancers get buried by 200 proposals. So I'm building the alternative. This is my first post BFS Algorithm in Java Step by Step Tutorial with Examples Tracking LLM Pricing Monthly: An Open Dataset for 22 AI Models How We Measure Content ROI on a Comparison Site: Revenue Attribution Without Perfect Data Introducing Nova AI Ops: The AI-Native Operating System for SRE Teams I built a free desktop video downloader for Windows — Grabbit How Talkie OCR Helps Vision-Impaired & Dyslexic Users Read the World Around Them VRCFaceTracking安装和iPhone面捕配置教程,有bug Even CrowdStrike Can't See Your Agents The Automation Gold Rush: What n8n Workflows and Claude Are Opening Up for Developers Right Now
Switching from file-based jobs to NATS/Kafka in Rust without changing code
Marco Mengel · 2026-04-17 · via DEV Community

Introduction

Most applications eventually need to offload work to another process. Parse a file, send an email, trigger a report – tasks that shouldn't block your main service and might even run on a different machine.

Most queue or stream solutions require you to commit to their ecosystem from day one – their broker, their worker format, their retry logic. And if you want to swap the transport later, due to a license change or some request from a good paying customer, you're rewriting business logic.

I wanted something simpler: define your jobs once in plain Rust structs, start with a file during development, and switch to a real broker for production – without touching the handler code.

This is how I built a remote job system in Rust using mq-bridge.

Step 1: Create Cargo.toml

Let's run cargo init, cargo add mq-bridge serde tokio tracing tracing-subscriber and some other modifications:

# Cargo.toml
[package]
name = "mq-bridge-jobs-example"
version = "0.1.0"
edition = "2021"

[dependencies]
mq-bridge = "0.2.11"
serde = { version = "1.0.228", features = ["derive"] }
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
tracing = "0.1.44"
tracing-subscriber = { version = "0.3.23", features = ["env-filter"] }

[[bin]]
name = "worker"
path = "src/bin/worker.rs"

[[bin]]
name = "submit"
path = "src/bin/submit.rs"

Enter fullscreen mode Exit fullscreen mode

We want 2 separate binaries: worker that waits for tasks and submit that sends a single mail, which should be received by our worker.


Step 2: Define your jobs

Before we touch any infrastructure, we define what our jobs look like. Just plain Rust structs:

// src/jobs.rs
use serde::{Serialize, Deserialize};

#[derive(Serialize, Deserialize)]
pub struct SendEmail {
    pub to: String,
    pub subject: String,
}

#[derive(Serialize, Deserialize)]
pub struct GenerateReport {
    pub user_id: u32,
}

Enter fullscreen mode Exit fullscreen mode

In addition, we define strings to identify each struct:

// src/jobs.rs
impl SendEmail {
    pub const KIND: &'static str = "send_email";
}

impl GenerateReport {
    pub const KIND: &'static str = "generate_report";
}

Enter fullscreen mode Exit fullscreen mode

Let's also add a lib.rs file:

// src/lib.rs
pub mod jobs;

Enter fullscreen mode Exit fullscreen mode

Then we register handlers for each job type using mq-bridge TypeHandler:

// src/bin/worker.rs
let jobs = TypeHandler::new()
    .add(SendEmail::KIND, |job: SendEmail| async move {
        // We are not actually sending a mail here - just print a log message
        tracing::info!("Sending email to {}", job.to);
        tokio::time::sleep(Duration::from_millis(100)).await;
        Ok(Handled::Ack)
    })
    .add(GenerateReport::KIND, |job: GenerateReport| async move {
        tracing::info!("Generating report for user {}", job.user_id);
        Ok(Handled::Ack)
    });

Enter fullscreen mode Exit fullscreen mode


Step 3: Start with a file backend

No Docker. No broker. Just a file on disk for our worker.

// src/bin/worker.rs
//...
let route = Route::new(
    Endpoint::new(EndpointType::File(
        FileConfig::new("jobs.jsonl").with_mode(FileConsumerMode::Consume { delete: true }),
    )),
    Endpoint::null(), // No output needed here
).with_handler(jobs);

route.deploy("job_worker").await?;

Enter fullscreen mode Exit fullscreen mode

Together with logging and everything, the complete worker.rs now looks like this:

// src/bin/worker.rs (complete)
use mq_bridge::{
    Handled, Route,
    models::{Endpoint, EndpointType, FileConfig, FileConsumerMode},
    type_handler::TypeHandler,
};
use mq_bridge_jobs_example::jobs::{GenerateReport, SendEmail};
use std::time::Duration;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    tracing_subscriber::fmt()
        .with_env_filter(tracing_subscriber::EnvFilter::new("info"))
        .init();

    let jobs = TypeHandler::new()
        .add(SendEmail::KIND, |job: SendEmail| async move {
            tracing::info!("Sending email to {}", job.to);
            tokio::time::sleep(Duration::from_millis(100)).await;
            Ok(Handled::Ack)
        })
        .add(GenerateReport::KIND, |job: GenerateReport| async move {
            tracing::info!("Generating report for user {}", job.user_id);
            Ok(Handled::Ack)
        });

    let route = Route::new(
        Endpoint::new(EndpointType::File(
            FileConfig::new("jobs.jsonl").with_mode(FileConsumerMode::Consume { delete: true }),
        )),
        Endpoint::null(), // No output needed here
    )
    .with_handler(jobs);

    route.deploy("job_worker").await?;

    tracing::info!("Worker running — press Ctrl-C to exit");
    tokio::signal::ctrl_c().await?;
    tracing::info!("Shutting down");
    Ok(())
}

Enter fullscreen mode Exit fullscreen mode

To submit a job, just append a new line to jobs.jsonl in our submit.rs:

// src/bin/submit.rs
use mq_bridge::{
    Publisher,
    models::{Endpoint, EndpointType, FileConfig},
    msg,
};
use mq_bridge_jobs_example::jobs::SendEmail;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    tracing_subscriber::fmt()
        .with_env_filter(tracing_subscriber::EnvFilter::new("info"))
        .init();

    let publisher = Publisher::new(Endpoint::new(EndpointType::File(FileConfig::new(
        "jobs.jsonl",
    ))))
    .await?;

    publisher
        .send(msg!(
            SendEmail {
                to: "user@example.com".into(),
                subject: "Welcome!".into()
            },
            SendEmail::KIND
        ))
        .await?;
    Ok(())
}

Enter fullscreen mode Exit fullscreen mode

Works completely offline. Great for development and testing.

Now let's test it. Open a first shell:

cargo run --bin worker

Enter fullscreen mode Exit fullscreen mode

The worker is now running and waiting for file modifications. In a second shell, submit a job:

cargo run --bin submit

Enter fullscreen mode Exit fullscreen mode

The worker will receive the task and print:

INFO worker: Sending email to user@example.com

Enter fullscreen mode Exit fullscreen mode

Instead of using the submit binary, you could also just simply push a new line to the file

echo '{"message_id":1,"payload":{"subject":"Welcome!","to":"user@example.com"},"metadata":{"kind":"send_mail"}}' > jobs.jsonl

Enter fullscreen mode Exit fullscreen mode

Afterwards jobs.jsonl is empty — because FileConsumerMode::Consume { delete: true } removes consumed lines. With delete: false, lines would be kept and replayed on the next worker start.

There is alternatively a GroupSubscribe mode to prevent re-deliver by tracking the current offset via separate .offset file, without deleting lines.


Step 4: Switch to JSON config

The business logic stays in Rust. The infrastructure moves to config:

cargo add serde_json

Enter fullscreen mode Exit fullscreen mode

src/bin/config.json

{
  "input": {
    "file": {
      "path": "jobs.jsonl",
      "delete": true,
      "mode": "consume"
    }
  },
  "output": {
    "null": null
  }
}

Enter fullscreen mode Exit fullscreen mode

// src/bin/worker.rs - load route from config
let route: Route = serde_json::from_str(include_str!("config.json"))?;
let route = route.with_handler(jobs);
route.deploy("job_worker").await?;

Enter fullscreen mode Exit fullscreen mode

// src/bin/submit.rs - create a publisher from the same config
let route: Route = serde_json::from_str(include_str!("config.json"))?;
let publisher = Publisher::new(route.input).await?;

Enter fullscreen mode Exit fullscreen mode

You can now load the configuration from a file or database. The code is smaller, and you can change the backend without touching your handler code.

In a later production scenario, you might also want to use a separate publisher configuration. The are properties that are only available for consumers or publishers and there would be a warning when using invalid settings. Also, you might want to configure a specific kafka group_id or use separate topics for fan out.
But for this example, using a common NATS configuration works fine.


Step 5: Switch to NATS for production

To run the worker on a separate machine, you'll want a broker or database. NATS is a great fit — it's lightweight, just a single binary with no dependencies and stores messages.

First, enable the nats feature in Cargo.toml:

mq-bridge = { version = "0.2.11", features = ["nats"] }

Enter fullscreen mode Exit fullscreen mode

This just enables the "nats" feature. We can simply re-run the previous
example. Nothing changes yet, still using file, it just needs longer to compile.

Start NATS with JetStream:

# macOS
brew install nats-server && nats-server -js

# or Ubuntu/Debian
wget https://github.com/nats-io/nats-server/releases/latest/download/nats-server-linux-amd64.deb
sudo apt install ./nats-server-linux-amd64.deb && nats-server -js

# or Docker
docker run -p 4222:4222 nats:2.12.2 -js

Enter fullscreen mode Exit fullscreen mode

One config.json file change, no code changes:

{
  "input": {
    "nats": {
      "url": "nats://localhost:4222",
      "subject": "test-stream.pipeline",
      "stream": "test-stream"
    }
  },
  "output": {
    "null": null
  }
}

Enter fullscreen mode Exit fullscreen mode

Restart worker and submit — both now talk to NATS. The handler code is untouched.


What you get for free

Switching to NATS unlocks everything mq-bridge builds on top. You can add middlewares in the config, for example retries and a dead-letter queue (DLQ) for failed messages:

{
  "nats": {
    "url": "nats://localhost:4222",
    "subject": "test-stream.pipeline",
    "stream": "test-stream"
  },
  "middlewares": [
    {
      "retry": {
        "max_attempts": 3,
        "max_interval_ms": 5000,
        "initial_interval_ms": 100,
        "multiplier": 2
      }
    },
    {
      "dlq": {
        "endpoint": {
          "file": {
            "path": "error.log"
          }
        }
      }
    }
  ]
}

Enter fullscreen mode Exit fullscreen mode

The retry middleware will retry failed deliveries with exponential backoff. If all attempts are exhausted, the dlq middleware writes the message to error.log instead of dropping it silently.

If you are already using MongoDB, MySQL, MariaDB, or PostgreSQL, you can use them as your queue backend as well — just a config change.

If you just want message forwarding from one endpoint to another or an UI to
create different json configs, you can also use
mq-bridge-app (cargo install mq-bridge-app).
The code also shows how you would use mq-bridge as webserver.

If you just need a simple send and receive - this is also available. You may skip the
whole event handler and route concept and just use the same API calls for Http, gRPC, MongoDb, Kafka, RabbitMQ and NATS. They all have the same receive and publish method and use the same message struct CanonicalMessage for transport.


Step 6: Testing with the memory endpoint

Because mq-bridge uses the same trait for all backends, you can test your handlers without any broker or file system — just an in-memory channel.

Let's add a test for submit.rs:


// src/bin/submit.rs
use mq_bridge::{msg, Publisher, Route};
use mq_bridge_jobs_example::jobs::SendEmail;

async fn send_mail(publisher: Publisher) -> Result<(), Box<dyn std::error::Error>> {
    publisher
        .send(msg!(
            SendEmail {
                to: "user@example.com".into(),
                subject: "Welcome!".into()
            },
            SendEmail::KIND
        ))
        .await?;
    Ok(())
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    tracing_subscriber::fmt()
        .with_env_filter(tracing_subscriber::EnvFilter::new("info"))
        .init();

    let route: Route = serde_json::from_str(include_str!("config.json"))?;
    let publisher = Publisher::new(route.input).await?;
    send_mail(publisher).await
}
#[cfg(test)]
mod tests {
    use mq_bridge::endpoints::memory::MemoryConsumer;
    use mq_bridge::traits::MessageConsumer;
    use mq_bridge::Publisher;
    use mq_bridge::models::Endpoint;
    use mq_bridge_jobs_example::jobs::SendEmail;

    use crate::send_mail;

    #[tokio::test]
    async fn test_submit_sends_email_job() {
        let topic = "test-submit";
        let mut consumer = MemoryConsumer::new_local(topic, 10);
        let publisher = Publisher::new(Endpoint::new_memory(topic, 10)).await.unwrap();
        send_mail(publisher).await.unwrap();
        let received = consumer.receive().await.unwrap();
        let payload: serde_json::Value = serde_json::from_slice(&received.message.payload).unwrap();
        assert_eq!(payload["to"], "user@example.com");
        assert_eq!(received.message.metadata["kind"], SendEmail::KIND);
    }
}

Enter fullscreen mode Exit fullscreen mode

No broker, no file, no test containers. The same TypeHandler that runs in production is tested here — only the transport is swapped.

What not to expect

Not all aspects and features of brokers or databases are supported. Some features
are emulated, other features may not be implemented yet. Don't expect a full grown
framework that guides you on how to do stuff or already prevents misconfiguration
during compile time when reading configs during runtime.


Conclusion

mq-bridge covers more than just remote jobs. You can use it for events, or to send and receive messages from existing brokers. And you can scale up by adding Kafka as a buffer or fan-out layer — again, just config.

mq-bridge is still a young library. Don't expect it to be as complete as Watermill (Go) or Java Spring. It uses some of their concepts, but it doesn't try to be the same — event sourcing and aggregate management are out of scope for now, as the focus is on transport. Documentation is still growing, and this tutorial is a first step toward that.

This tutorial is available here:
https://github.com/marcomq/mq-bridge-jobs-example

The mq-bridge library is available here:
(I’m the author of mq-bridge, just for transparency.)
https://github.com/marcomq/mq-bridge

Feedback and contributions welcome.