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

推荐订阅源

AWS News Blog
AWS News Blog
T
Tenable Blog
Project Zero
Project Zero
T
The Exploit Database - CXSecurity.com
L
LINUX DO - 热门话题
T
Threat Research - Cisco Blogs
T
Threatpost
Security Latest
Security Latest
C
Cisco Blogs
L
Lohrmann on Cybersecurity
S
Security @ Cisco Blogs
Google Online Security Blog
Google Online Security Blog
NISL@THU
NISL@THU
AI
AI
V
Vulnerabilities – Threatpost
Google DeepMind News
Google DeepMind News
C
Cyber Attacks, Cyber Crime and Cyber Security
C
CXSECURITY Database RSS Feed - CXSecurity.com
The Last Watchdog
The Last Watchdog
G
GRAHAM CLULEY
Cloudbric
Cloudbric
H
Hackread – Cybersecurity News, Data Breaches, AI and More
H
Hacker News: Front Page
U
Unit 42
A
Arctic Wolf
CTFtime.org: upcoming CTF events
CTFtime.org: upcoming CTF events
MyScale Blog
MyScale Blog
O
OpenAI News
Scott Helme
Scott Helme
V2EX - 技术
V2EX - 技术
P
Proofpoint News Feed
博客园 - 叶小钗
Hugging Face - Blog
Hugging Face - Blog
云风的 BLOG
云风的 BLOG
V
Visual Studio Blog
Application and Cybersecurity Blog
Application and Cybersecurity Blog
Cyberwarzone
Cyberwarzone
博客园 - 【当耐特】
H
Heimdal Security Blog
S
Schneier on Security
阮一峰的网络日志
阮一峰的网络日志
Help Net Security
Help Net Security
D
DataBreaches.Net
Y
Y Combinator Blog
Hacker News - Newest:
Hacker News - Newest: "LLM"
TaoSecurity Blog
TaoSecurity Blog
K
Kaspersky official blog
N
News and Events Feed by Topic
WordPress大学
WordPress大学
P
Palo Alto Networks 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 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
Designing a Scalable Event-Driven Data Processing Pipeline with Apache Kafka Streams
Rizwan Saleem · 2026-05-31 · via DEV Community

Rizwan Saleem

Designing a Scalable Event-Driven Data Processing Pipeline with Apache Kafka Streams

Designing a Scalable Event-Driven Data Processing Pipeline with Apache Kafka Streams

In modern data-intensive applications, real-time insights often drive user value. A robust event-driven data processing pipeline lets you ingest, transform, and route data with low latency while remaining resilient to failures and traffic bursts. This guide walks through designing and implementing a scalable, maintainable event-driven pipeline using Apache Kafka and Kafka Streams. It covers architecture decisions, data modeling, fault tolerance, deployment, and practical code examples you can adapt to your stack.

Overview of the architecture

  • Event producer layer: services that emit events in well-defined schemas.
  • Event broker: Apache Kafka clusters that persist events and decouple producers from consumers.
  • Stream processing layer: Kafka Streams applications that transform, enrich, and route data in real time.
  • Sinks and consumers: downstream databases, caches, search indices, or microservices that react to processed results.
  • Operational tooling: monitoring, schema management, deployments, and testing.

Key design principles

  • Stateless stream processing: keep processors idempotent and stateless where possible to simplify scaling and recovery.
  • Exactly-once semantics (EOS) where needed: configure Kafka and streams to minimize duplicate processing in critical paths.
  • Loose coupling via schemas: use a strong schema on read/write to evolve data safely.
  • Backpressure-aware design: handle backpressure gracefully to avoid data loss or unbounded buffering.
  • Observability by design: instrument metrics, traces, and logs at producers, streams, and sinks.

1) Data modeling and schemas

  • Choose a canonical event schema: define clear event types (e.g., UserCreated, OrderPlaced, InventoryUpdated) with common envelope fields:
    • event_id, timestamp, source, type, payload.
  • Use a schema registry (e.g., Confluent Schema Registry) to manage Avro/JSON schemas and enforce compatibility.
  • Version schemas: avoid breaking changes by versioning events or introducing new event types without altering existing ones.
  • Idempotent payloads: design payloads so reprocessing the same event yields the same result (upserts, up-to-date views).

Example Avro payload for an OrderPlaced event:

  • { "type": "record", "name": "OrderPlaced", "fields": [ {"name": "order_id", "type": "string"}, {"name": "user_id", "type": "string"}, {"name": "items", "type": {"type": "array", "items": { "type": "record", "name": "Item", "fields": [ {"name": "sku", "type": "string"}, {"name": "quantity", "type": "int"}, {"name": "price", "type": "double"} ] }}}, {"name": "total", "type": "double"}, {"name": "timestamp", "type": "long"} ] }

2) Ingest layer: producers and topics

  • Organize topics by event type and domain boundaries (e.g., orders-raw, orders-processed, inventory-events).
  • Partitioning strategy: partition by a key that ensures related events land on the same partition to improve locality (e.g., user_id for user-centric events, order_id for order events).
  • At-least-once vs exactly-once delivery:
    • Producers: enable idempotence and acks=all.
    • For EOS: enable transactional producers if you have multi-topic writes in a single unit of work.
  • Serialization:
    • Use Avro with schema registry for compact binary payloads and strong typing.
    • Ensure producer and consumer clients share the same schema.

Code sketch: Kafka producer with Avro and transactional support (Java/Spring Boot style)

  • Note: this is a simplified sketch; integrate with your framework and error handling.

  • Dependencies: org.apache.kafka:kafka-clients, io.confluent:kafka-avro-serializer

  • Pseudo-implementation:

    • properties:
    • bootstrap.servers
    • key.serializer=org.apache.kafka.common.serialization.StringSerializer
    • value.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
    • enable.idempotence=true
    • transactional.id=order-producer-1
    • producer.initTransactions()
    • producer.beginTransaction()
    • producer.send(new ProducerRecord<>("orders-raw", key, orderEvent))
    • producer.commitTransaction()

3) Stream processing layer: Kafka Streams

  • Use a topology that handles:
    • Filtering and enriching: join with reference data, compute derived metrics.
    • Windowed aggregations for real-time dashboards (e.g., 1-minute tumbling windows).
    • Outbox pattern for exactly-once guarantees in downstream sinks.
  • State stores:
    • Use RocksDB-backed stores for local state; tune cache size and segment purge settings.
    • Use changelog topics to restore state on restart.
  • Fault tolerance:
    • Enable EOS where required, ensure consumer groups have proper isolation.
    • Configure commit interval and processing guarantees to balance latency and throughput.
  • Scalability:
    • Increase num.partitions for topics with high throughput.
    • Run multiple instances of the Streams app; Kafka partitions guide parallelism.
  • Observability:
    • Emit metrics for throughput, lag, processing time, and error rates.
    • Use tracing (OpenTelemetry) to trace end-to-end flow across producers, streams, and sinks.

Code sketch: Kafka Streams topology (Java)

  • Build a topology that consumes orders-raw, enriches with customer data from a KTable, and writes to orders-processed and an outbox topic.

  • Dependencies: org.apache.kafka:kafka-streams

  • Pseudo-implementation:

    • StreamsBuilder builder = new StreamsBuilder();
    • KStream orders = builder.stream("orders-raw", Consumed.with(string(), avro(OrderEvent.class)));
    • KTable customers = builder.table("customers", Consumed.with(string(), avro(CustomerInfo.class)));
    • KStream enriched = orders.join(customers, (order, customer) -> enrichOrder(order, customer), Joined.with(string(), avro(OrderEvent.class), avro(CustomerInfo.class)) );
    • enriched through a windowed aggregation or mapping to processed structure.
    • enriched.to("orders-processed", Produced.with(string(), avro(ProcessedOrder.class)));
    • enriched.mapValues(p -> p.toOutbox()).to("orders-outbox", Produced.with(string(), avro(OutboxEvent.class)));
    • KafkaStreams streams = new KafkaStreams(build(), config);
    • streams.start();

4) Sinks and downstream consumers

  • Processed data sink: databases (PostgreSQL, Cassandra), search indexes (Elasticsearch), or caches (Redis).
  • Outbox pattern:
    • Produce events to an outbox topic as a reliable bridge to sinks.
    • A separate consumer reads outbox events and writes to the sink, with idempotent writes to prevent duplicates.
  • Error handling:
    • Dead-letter queues (DLQ) for failed messages with metadata to diagnose issues.
  • Idempotent sinks:
    • Upsert semantics or resource-based locks to avoid duplicates.

Example: Outbox consumer sketch (Java)

  • Consumes from orders-outbox and writes to a relational database using upsert.
  • On failure, publishes to a DLQ topic with error metadata.

5) Deployment and operations

  • Deployments:
    • Separate deploys for producers, streams, and sinks to minimize blast radius.
    • Use containerization (Docker) or serverless-like environments (Kubernetes) with resource requests/limits.
  • Configuration management:
    • Externalize config with environment variables or a config service.
    • Versioned configurations to track changes over time.
  • Observability:
    • Metrics: throughput (records/sec), latency (end-to-end), lag (consumer group lag), error counts.
    • Tracing: propagate trace context across producers and streams for end-to-end visibility.
    • Logs: structured logs with correlation IDs.
  • Resilience and failure recovery:
    • Replication factor and in-sync replica settings on Kafka topics.
    • Enable topic-level compaction where appropriate to prune old data.
    • Regular backup and restore drills for Kafka clusters and sinks.

6) Testing strategies

  • Unit tests:
    • Mock Kafka topics and test topologies with TopologyTestDriver.
  • Integration tests:
    • End-to-end tests using a test containerized Kafka cluster (e.g., Testcontainers) and a lightweight sink.
  • Chaos testing:
    • Simulate broker outages, network partitions, and lag to observe recovery behavior.
  • Data quality checks:
    • Validate schema compatibility, field presence, and required values.

7) Practical, end-to-end example walk-through

Scenario: Real-time order analytics dashboard

  • Producers emit:
    • orders-raw: OrderPlaced events with order_id, user_id, items, total, timestamp.
    • users-raw: UserCreated events for enriching customer data.
  • Streams:
    • Consume orders-raw and users-raw; enrich orders with user segment from a KTable built from users-raw.
    • Compute per-minute revenue by region and user segment, store in orders-processed.
    • Emit outbox events for dashboards and anomaly detection.
  • Sinks:
    • orders-processed writes to PostgreSQL for dashboards.
    • orders-outbox consumed by a separate service pushing metrics to a real-time dashboard (e.g., Grafana) via time-series store.

Code snippet: End-to-end data flow in outline

  • Producer writes to orders-raw with transactional producer to ensure single-unit writes.
  • Streams topology reads orders-raw and users, joins them, writes to orders-processed and orders-outbox.
  • Outbox consumer updates dashboards and alerting systems.

8) Security considerations

  • Encrypt data in transit:
    • Enable TLS for all Kafka clients.
  • Data at rest:
    • Enable disk encryption on cluster storage if required.
  • Access control:
    • Use SCRAM or OAuth for authentication; apply ACLs per topic.
  • Data governance:
    • Mask or redact sensitive fields in non-secure paths; consider separate topics for sensitive vs non-sensitive data.

9) Common pitfalls and tips

  • Latency vs throughput trade-offs:
    • Lower commit intervals reduce latency but increase commit overhead; tune accordingly.
  • Schema evolution:
    • Prefer additive changes and backward compatibility; avoid breaking changes in live streams.
  • Backpressure:
    • Implement buffering limits and backpressure-aware sinks; monitor consumer lag.
  • Operational complexity:
    • Start small with a minimal viable pipeline and gradually add enrichments and sinks.

If you want, I can tailor this guide to a specific tech stack (e.g., Java/Kotlin, Python with Faust, or Node.js with KafkaJS) or adapt it to a particular domain (e-commerce, IoT telemetry, or financial data). Would you like a concrete, language-specific code example continued for your preferred stack?

-

Rizwan Saleem | https://rizwansaleem.co