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

推荐订阅源

H
Help Net Security
J
Java Code Geeks
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
H
Hackread – Cybersecurity News, Data Breaches, AI and More
V
Visual Studio Blog
G
Google Developers Blog
V
V2EX
The Register - Security
The Register - Security
博客园 - 三生石上(FineUI控件)
云风的 BLOG
云风的 BLOG
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
CTFtime.org: upcoming CTF events
CTFtime.org: upcoming CTF events
博客园_首页
S
SegmentFault 最新的问题
博客园 - Franky
Martin Fowler
Martin Fowler
Stack Overflow Blog
Stack Overflow Blog
A
About on SuperTechFans
人人都是产品经理
人人都是产品经理
aimingoo的专栏
aimingoo的专栏
罗磊的独立博客
C
Check Point Blog
MyScale Blog
MyScale Blog
T
The Blog of Author Tim Ferriss
MongoDB | Blog
MongoDB | Blog
The GitHub Blog
The GitHub Blog
Last Week in AI
Last Week in AI
Microsoft Azure Blog
Microsoft Azure Blog
IT之家
IT之家
F
Fortinet All Blogs
Jina AI
Jina AI
P
Proofpoint News Feed
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
阮一峰的网络日志
阮一峰的网络日志
B
Blog
L
LangChain Blog
月光博客
月光博客
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
宝玉的分享
宝玉的分享
博客园 - 【当耐特】
T
Tailwind CSS Blog
酷 壳 – CoolShell
酷 壳 – CoolShell
Microsoft Security Blog
Microsoft Security Blog
WordPress大学
WordPress大学
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
B
Blog RSS Feed
博客园 - 聂微东
Hugging Face - Blog
Hugging Face - Blog
M
MIT News - Artificial intelligence
GbyAI
GbyAI

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 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 Switching from file-based jobs to NATS/Kafka in Rust without changing code 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
I Built a Mini Message Broker in Pure Python and Finally Understood How Kafka Moves Millions of Events
Haji Rufai · 2026-06-16 · via DEV Community

Last year I was on a team that pushed 40 million events per day through Kafka. We had consumer lag alerts, rebalancing incidents, and a whole runbook for when the broker got behind. I understood how to operate Kafka. But I did not understand how Kafka works.

So I built a tiny one. No dependencies. No Zookeeper. No JVM. Just Python and the core ideas.

Here is what I learned.


The Three Things Kafka Actually Does

People say "Kafka is a message queue." That is not quite right.

Kafka is a distributed commit log. It has three jobs:

  1. Accept writes from producers and append them to a log
  2. Let consumers read from any offset in that log
  3. Remember where each consumer group is up to

That third one is the thing that makes Kafka different from a traditional queue. A queue forgets a message once it is consumed. Kafka remembers. You can replay. You can have 10 different consumer groups reading the same topic at different speeds.

The code to implement this is smaller than you think.


brokelite: A Message Broker in 120 Lines

import threading
import time
from collections import defaultdict
from typing import Dict, List, Tuple

class Partition:
    """Append-only log for one partition of a topic."""

    def __init__(self):
        self._log: List[Tuple[int, bytes]] = []  # (offset, message)
        self._lock = threading.Lock()
        self._next_offset = 0

    def append(self, message: bytes) -> int:
        with self._lock:
            offset = self._next_offset
            self._log.append((offset, message))
            self._next_offset += 1
            return offset

    def read_from(self, offset: int, max_count: int = 100) -> List[Tuple[int, bytes]]:
        with self._lock:
            return [
                (off, msg)
                for off, msg in self._log
                if off >= offset
            ][:max_count]

    def __len__(self):
        return self._next_offset


class Topic:
    """A topic is just N partitions."""

    def __init__(self, name: str, num_partitions: int = 3):
        self.name = name
        self.partitions = [Partition() for _ in range(num_partitions)]

    def route(self, key: bytes | None) -> int:
        """Route to a partition by key hash, or round-robin if no key."""
        if key is None:
            return int(time.time() * 1000) % len(self.partitions)
        return hash(key) % len(self.partitions)

    def produce(self, message: bytes, key: bytes | None = None) -> Tuple[int, int]:
        partition_id = self.route(key)
        offset = self.partitions[partition_id].append(message)
        return partition_id, offset

    def consume(self, partition_id: int, offset: int, max_count: int = 100):
        return self.partitions[partition_id].read_from(offset, max_count)


class ConsumerGroupRegistry:
    """Tracks committed offsets per group, topic, partition."""

    def __init__(self):
        # {group: {topic: {partition: offset}}}
        self._offsets: Dict[str, Dict[str, Dict[int, int]]] = defaultdict(
            lambda: defaultdict(lambda: defaultdict(int))
        )
        self._lock = threading.Lock()

    def commit(self, group: str, topic: str, partition: int, offset: int):
        with self._lock:
            self._offsets[group][topic][partition] = offset

    def get_offset(self, group: str, topic: str, partition: int) -> int:
        with self._lock:
            return self._offsets[group][topic][partition]


class Broker:
    """The broker: accepts produces, serves consumes, tracks consumer groups."""

    def __init__(self):
        self._topics: Dict[str, Topic] = {}
        self._registry = ConsumerGroupRegistry()
        self._lock = threading.Lock()

    def create_topic(self, name: str, num_partitions: int = 3):
        with self._lock:
            if name not in self._topics:
                self._topics[name] = Topic(name, num_partitions)

    def produce(self, topic: str, message: bytes, key: bytes | None = None):
        if topic not in self._topics:
            self.create_topic(topic)
        partition_id, offset = self._topics[topic].produce(message, key)
        return {"topic": topic, "partition": partition_id, "offset": offset}

    def consume(
        self, topic: str, group: str, partition: int, max_count: int = 100
    ) -> List[Tuple[int, bytes]]:
        if topic not in self._topics:
            return []
        offset = self._registry.get_offset(group, topic, partition)
        records = self._topics[topic].consume(partition, offset, max_count)
        return records

    def commit(self, topic: str, group: str, partition: int, offset: int):
        self._registry.commit(group, topic, partition, offset)

That is the whole broker. Let me walk through the key ideas.


The Partition is an Append-Only Log

def append(self, message: bytes) -> int:
    with self._lock:
        offset = self._next_offset
        self._log.append((offset, message))
        self._next_offset += 1
        return offset

Every message gets a monotonically increasing offset. That offset never changes. That is the core promise. A Kafka offset is like a line number in a file. Once written, line 42 is always line 42.

This is why Kafka is fast. Appends are O(1). You never update old records. There is no random I/O.

Real Kafka writes these to segment files on disk in sequential order. Sequential writes on spinning disk are nearly as fast as sequential writes on SSD. That is why Kafka can hit hundreds of MB/s throughput on commodity hardware.


Key-Based Routing Gives You Ordering Guarantees

def route(self, key: bytes | None) -> int:
    if key is None:
        return int(time.time() * 1000) % len(self.partitions)
    return hash(key) % len(self.partitions)

Here is a thing that trips up new Kafka users: Kafka only guarantees ordering within a partition. Not across partitions.

If you produce events for user_id=123 and they land on partition 0, 2, and 1 in that order, your consumer sees them out of sequence.

The fix is simple: use a consistent key. hash(user_id) % num_partitions means all events for user 123 always land on the same partition. Ordering preserved.

When you do not provide a key, the broker round-robins. Good for throughput. Bad for ordering. Know which one you need.


Consumer Groups Enable Independent Progress

def get_offset(self, group: str, topic: str, partition: int) -> int:
    with self._lock:
        return self._offsets[group][topic][partition]

This is the idea that changed how I think about event-driven systems.

In a traditional queue, there is one cursor. One consumer group owns it. If you want a second application to process the same events, you have to publish to two separate queues. Or copy the data.

Kafka flips this. The offset belongs to the consumer group, not the broker. Every group maintains its own pointer. Group A can be at offset 1000. Group B can be at offset 500. The broker does not care. It keeps the log until retention expires.

This is why Kafka works so well for fan-out. Your analytics team can replay from offset 0 without blocking your real-time alerting pipeline.


Using brokelite

broker = Broker()

# Producer side
broker.produce("orders", b'{"order_id": 1, "amount": 99.99}', key=b"user_123")
broker.produce("orders", b'{"order_id": 2, "amount": 14.50}', key=b"user_456")
broker.produce("orders", b'{"order_id": 3, "amount": 7.25}', key=b"user_123")

# Consumer side: group A reads partition 0
records = broker.consume("orders", group="analytics", partition=0, max_count=10)
for offset, msg in records:
    print(f"offset={offset}: {msg.decode()}")
    broker.commit("orders", "analytics", partition=0, offset=offset + 1)

# Consumer side: group B reads the same partition independently
records = broker.consume("orders", group="alerting", partition=0, max_count=10)
for offset, msg in records:
    print(f"alerting sees offset={offset}: {msg.decode()}")

Two consumer groups, same data, completely independent offsets. Zero copying. This is the elegance of the commit log model.


What brokelite Leaves Out

This is a toy. Real Kafka has things this does not.

Replication. Each partition in real Kafka has a leader and N replicas. Writes go to the leader and replicate before being acknowledged (depending on acks setting). brokelite has no replication and no durability beyond memory.

Consumer group rebalancing. In real Kafka, when a consumer in a group dies, the partitions it owned get reassigned to surviving consumers. That is the source of most Kafka operational pain. brokelite has no concept of group membership.

Log compaction and retention. Real Kafka can either delete old segments after N days or compact them by key. This keeps disk usage bounded. brokelite grows forever.

Network protocol. Real Kafka has a binary TCP protocol. Producers and consumers connect over the network. brokelite is in-process only.

Those are real engineering problems. But once you understand the core, those features are additive. They do not change the fundamental model: append-only log, per-group offsets, key-based partitioning.


The Lesson I Keep Coming Back To

Every time an alert fires for consumer lag, I used to think "the broker is behind." That is wrong.

Consumer lag is log_end_offset - committed_offset for a consumer group. The broker is not behind. The consumer is behind. The broker is fine, it is just sitting there with messages no one has read yet.

This distinction matters when you are debugging. Slow consumer group? Check your consumer throughput, not your broker. Skewed lag across partitions? Check your key distribution.

Build the toy version. Read the log. The model becomes obvious.


What to Build Next

If you want to go further with brokelite:

  • Add a seek_to_beginning method so a group can replay from offset 0
  • Add a seek_to_end method so a new group skips history and reads only new messages
  • Add a simple retention policy that purges segments older than N seconds
  • Simulate a consumer dying mid-batch and see what happens to uncommitted offsets

The more you push on the toy, the better your intuition gets for the real system.

Follow along for more deep dives into systems that data engineers use every day but rarely look inside.

What pattern do you rely on most in your streaming pipelines: key-based ordering or maximum throughput with no keys?