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

推荐订阅源

N
News and Events Feed by Topic
爱范儿
爱范儿
Blog — PlanetScale
Blog — PlanetScale
The GitHub Blog
The GitHub Blog
C
Check Point Blog
小众软件
小众软件
I
InfoQ
罗磊的独立博客
H
Hackread – Cybersecurity News, Data Breaches, AI and More
Engineering at Meta
Engineering at Meta
酷 壳 – CoolShell
酷 壳 – CoolShell
CTFtime.org: upcoming CTF events
CTFtime.org: upcoming CTF events
Hugging Face - Blog
Hugging Face - Blog
博客园 - 三生石上(FineUI控件)
MyScale Blog
MyScale Blog
The Cloudflare Blog
Last Week in AI
Last Week in AI
腾讯CDC
Y
Y Combinator Blog
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
雷峰网
雷峰网
B
Blog
T
Tailwind CSS Blog
MongoDB | Blog
MongoDB | Blog
A
About on SuperTechFans
D
Docker
博客园 - 司徒正美
博客园_首页
Recent Announcements
Recent Announcements
D
DataBreaches.Net
阮一峰的网络日志
阮一峰的网络日志
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
G
Google Developers Blog
Microsoft Security Blog
Microsoft Security Blog
F
Fortinet All Blogs
Stack Overflow Blog
Stack Overflow Blog
aimingoo的专栏
aimingoo的专栏
N
Netflix TechBlog - Medium
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
博客园 - 聂微东
GbyAI
GbyAI
Jina AI
Jina AI
V
V2EX
Vercel News
Vercel News
IT之家
IT之家
WordPress大学
WordPress大学
M
MIT News - Artificial intelligence
NISL@THU
NISL@THU
V
Visual Studio Blog
C
Cybersecurity and Infrastructure Security Agency CISA

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
Building a Real-Time Kafka + Cassandra Pipeline
GeraldM · 2026-05-24 · via DEV Community

Introduction

Apache Kafka and Apache Cassandra pair effectively because they complement each other's strengths: Kafka handles high throughput, real-time event streaming and ingestion, while Cassandra provides scalable, fault tolerant and low-latency persistent storage for processed data.

Example: A movies streaming company from their platform may be streaming billions of events per day including user viewing behavior, playback metrics and content recommendations. Kafka enables real-time streaming and processing of this events. These high velocity streams are then consumed and persisted into Cassandra that acts as a highly scalable, fault tolerant database for storing time-series data and user activity logs. With this combination, the movie company is able to achieve massive write throughput, low latency reads by recommendation engines and reliable handling of global traffic while maintaining high reliability. That is how Netflix does it.

What is Apache Cassandra?

It is a free, open-source NoSQL database designed to handle large volumes of data across multiple nodes using a columnar storage architecture. It supports both read and write operations one every node (a node is a single server or machine within the Cassandra cluster that stores data and handles read and write requests) enabling data replication across nodes and ensuring high availability without a single point of failure.

Setting up Apache Cassandra

The following steps show how to download and start Cassandra:

  1. Make a folder for Cassandra.

    mkdir cassandra
    


    shell

  2. Download Cassandra using the following command.

    wget https://dlcdn.apache.org/cassandra/5.0.8/apache-cassandra-5.0.8-bin.tar.gz
    
  3. Extract archive to Cassandra folder.

    tar -xvf apache-cassandra-5.0.8-bin.tar.gz
    


    shell

  4. Set up environment variables.

    export CASSANDRA_HOME=~/cassandra/apache-cassandra-5.0.8-bin
    export PATH=$PATH:$CASSANDRA_HOME/bin
    
  5. Start Cassandra by navigating into the created cassandra folder after extraction cd apache-cassandra-5.0.8 then run the command bin/cassandra.

  6. Start the CQL shell
     Note: Cassandra Query Language (CQL) is the primary query language for Apache Cassandra, designed to feel familiar like SQL while working with Cassandra’s distributed wide-column data model. Unlike traditional SQL, CQL operates on Keyspaces (databases), tables (wide-column structures) and supports partition keys and clustering columns for data distribution and on-disk sorting. It allows you to create tables, insert data, perform queries with SELECT, ** WHERE** and ORDER BY and use lightweight transactions.

  7. Create a Keyspace(Database) and set the keyspace to be in use.
    Using:

    CREATE KEYSPACE weather_data2 WITH REPLICATION = { 'class' : 'SimpleStrategy', 'replication_factor' : 1};
    


  8. Create a table in the Keyspace(database) and use a SELECT statement to confirm table creation.

Errors encountered:

  • Cassandra process being killed during startup: When this happens, it's likely because of not having enough memory. This can be addressed by editing the jvm-clients.options file and adding the following.
     This sets the amount of memory (ram) you want Cassandra to use. Initially, Cassandra calculates this values to half the total memory of your compute.

  • Not recommended to run Cassandra as root
     It is not an error but a recommendation. Avoid running Cassandra as root as there is a possibility of running into errors.
     Using bin/cassandra, a successful Cassandra start looks like this. Leave that terminal open and start another terminal to proceed with the next steps.
    To stop Cassandra use: pkill -f cassandra

Kafka-Cassandra Pipeline

Let's use an example of where we stream real-time weather data from OpenWeather API using Kafka and store it on Cassandra database.

1. Kafka Producer + OpenWeather API

from kafka import KafkaProducer
import requests
from dotenv import load_dotenv
import os
import json

load_dotenv()

API_KEY = os.getenv('API_KEY')

def pull_weather_data():
    cities = ['New York', 'London', 'Johannesburg', 'Nairobi', 'Cairo', 'Doha', 'Tokyo', 'Sydney']
    cities_weather_data = []

    for city in cities:
        url = f"https://api.openweathermap.org/data/2.5/weather?q={city}&appid={API_KEY}"
        response = requests.get(url)
        weather_data = response.json()

        cities_weather_data.append({
            'City': weather_data['name'],
            'Country': weather_data['sys']['country'],
            'Temparature': weather_data['main']['temp'],
            'Humidity': weather_data['main']['humidity'],
            'Feels_Like': weather_data['main']['feels_like'],
            'Last_update_time': weather_data['dt']
        })

    return cities_weather_data

producer = KafkaProducer(
    bootstrap_servers='localhost:9092',
    value_serializer=lambda v: json.dumps(v).encode('utf-8')
)

topic = 'open_weather_api_cities_data'

while True:
    weather_data = pull_weather_data()
    producer.send(topic, weather_data)
    print(f'From openweather: {weather_data}')
    time.sleep(10)

Enter fullscreen mode Exit fullscreen mode

The above python code (producer), utilizes a REST API to collect weather data of specified cities from OpenWeather and write the data into a Kafka topic.
To understand more on Kafka and Kafka producers, checkout the following article:
A Beginners guide to Real-time Data Streaming with Apache Kafka

  • Output you get after running the Python Producer code:

2. Kafka Consumer + Insert into Cassandra

from kafka import KafkaConsumer
from cassandra.cluster import Cluster
from cassandra.query import SimpleStatement
from datetime import datetime
import json

cluster = Cluster(['localhost'])
session = cluster.connect()
session.set_keyspace('kafka_data')

print('Connected to cassandra')

insert_query = SimpleStatement(
    '''
    INSERT INTO cities_weather_data(
        city,
        country,
        last_update_time,
        temperature,
        feels_like,
        humidity
    )
    values (%s, %s, %s, %s, %s, %s)
    '''
)

consumer = KafkaConsumer(
    'open_weather_api_cities_data',
    bootstrap_servers='localhost:9092',
    auto_offset_reset='earliest',
    value_deserializer=lambda m: json.loads(m.decode('utf-8'))
)

print('Connected to Kafka and Consumer started...')

for message in consumer:
    weather_dict_data = message.value
    print(weather_dict_data)

    for data in weather_dict_data:
        timestamp = datetime.fromtimestamp(data["Last_update_time"])

        session.execute(
            insert_query,
            (
                data['City'],
                data['Country'],
                timestamp,
                data['Temparature'],
                data['Feels_Like'],
                data['Humidity']
            )
        )

    print(f"Inserted: {data['City']}")

Enter fullscreen mode Exit fullscreen mode

  • Start by importing dependencies which include KafkaConsumer for connecting to Kafka and consume messages from a topic, Cluster which connects python to a Cassandra database cluster, SimpleStatement which prepares and executes Cassandra Query Language (CQL) statements, datetime to convert timestamps to python datetime and json to deserialize incoming Kafka json messages into python dictionaries.
  • Connect to Cassandra. Cluster(['localhost']) connects to a local instance of Cassandra, cluster.connect() creates a session used for communication with Cassandra and session.set_keyspace('kafka_data') instructs Cassandra to use the "kafka_data" database for all operations.
  • Prepare an insert Query using SimpleStatement() and store it in a variable.
  • Connect to Kafka. KafkaConsumer() connects to Kafka, subscribes to the defined topic "open_weather_api_cities_data", connects to the local running kafka instance using bootstrap_servers = 'localhost:9092', auto_offset_reset = 'earliest' instructs the consumer to start reading from the first message in the topic and then deserialize the incoming messages to python dictionaries.
  • Using an infinite for loop, read every message on the topic and keep the consumer open waiting for new incoming messages.
  • weather_dict_data = message.value extract the weather data from he messages and store it in a variable.
  • The producer sends a list of all the city weather records, using a for data in weather_dict_data loop, process one city at a time and insert the data into Cassandra using session.execute() which executes the prepared Cassandra query while appending the obtained data.
  • On our data from OpenWeather, our timestamp was in unix, using timestamp = datetime.fromtimestamp(data["Last_update_time"]) we convert it to datetime. From 1725344440 to 2024-09-03 10:20:40 Why? For human readability.
  • Lastly we print print(f"Inserted: {data['City']}") as successful insert message to confirm data insertion.
  • The output we get after running the python Consumer code:
  • Query the database to get the result of the written data on Cassandra

Conclusion

By consuming weather data from a Kafka topic, transforming it into a structured format and writing it into a Cassandra keyspace(Database), we have built a simple scalable architecture capable of handling continuous high-throughput data streams. This is pattern is a common foundation in modern data engineering systems where Kafka acts as the collection/ingestion layer and Cassandra database serves as a storage layer optimized for fast read and writes. Together, these technologies enable reliable end-to-end pipelines used in applications such as monitoring systems, IoT platforms and real-time analytics engines. From here, the architecture can be extended further to include processing layers such as Spark or Flink or routing data into search and observability tools like Elasticsearch.