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

推荐订阅源

MyScale Blog
MyScale Blog
量子位
宝玉的分享
宝玉的分享
爱范儿
爱范儿
云风的 BLOG
云风的 BLOG
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
Recent Announcements
Recent Announcements
Apple Machine Learning Research
Apple Machine Learning Research
N
News and Events Feed by Topic
TaoSecurity Blog
TaoSecurity Blog
博客园 - 三生石上(FineUI控件)
小众软件
小众软件
Simon Willison's Weblog
Simon Willison's Weblog
Google DeepMind News
Google DeepMind News
K
KPMG report finds enterprise disconnect between AI and its ROI | CIO
aimingoo的专栏
aimingoo的专栏
Cloudbric
Cloudbric
Blog — PlanetScale
Blog — PlanetScale
Latest news
Latest news
S
Security @ Cisco Blogs
Last Week in AI
Last Week in AI
cs.CV updates on arXiv.org
cs.CV updates on arXiv.org
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
Vercel News
Vercel News
W
WeLiveSecurity
M
MIT News - Artificial intelligence
P
Proofpoint News Feed
P
Proofpoint News Feed
P
Palo Alto Networks Blog
www.infosecurity-magazine.com
www.infosecurity-magazine.com
T
The Blog of Author Tim Ferriss
腾讯CDC
大猫的无限游戏
大猫的无限游戏
Martin Fowler
Martin Fowler
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
V
V2EX
H
Hackread – Cybersecurity News, Data Breaches, AI and More
Stack Overflow Blog
Stack Overflow Blog
IT之家
IT之家
有赞技术团队
有赞技术团队
Microsoft Security Blog
Microsoft Security Blog
Exploit-DB.com RSS Feed
Exploit-DB.com RSS Feed
美团技术团队
博客园 - 【当耐特】
D
DataBreaches.Net
I
InfoQ
G
GRAHAM CLULEY
S
SegmentFault 最新的问题
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
B
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
Building a Real-Time Weather Streaming Pipeline with Kafka, Docker & Python
Damaa-C · 2026-05-18 · via DEV Community

Introduction

In modern data engineering, handling high-velocity, real-time streams requires decoupled architectures that can scale seamlessly. A simple script fetching data from an API and writing it straight to a database creates a tight coupling; if the database goes down or the API experiences a spike, the entire system breaks.

This project implements a resilient Event-Driven Architecture (EDA). It extracts live global weather metrics from an API, streams them into an Apache Kafka topic managed via Docker, and processes them through an ETL (Extract, Transform, Load) consumer engine that flattens and persists the data into a PostgreSQL database.


Project System Architecture & Directory Layout

The project decouples data sourcing from data transformation and consumption using a publish-subscribe model. Docker isolates the streaming platform infrastructure, while Python applications drive the data operations.

text
openweather-kafka_confluent-project/
├── docker-compose.yml       # Orchestrates Zookeeper, Kafka Broker, & Control Center
├── producer.py              # Extracted RapidAPI multi-city pipeline (Ingestion)
├── consumer.py              # Advanced Pandas & SQLAlchemy Postgres pipeline (ETL)
├── test.ipynb               # Jupyter Notebook for interactive validation & debugging
└── .env                     # Local container and API credential configurations

Enter fullscreen mode Exit fullscreen mode

Infrastructure Layer: Docker Compose & Commands

Instead of dealing with local, environment-specific installations of Kafka, the entire messaging backbone is containerized. The docker-compose.yml provisions a robust Confluent platform stack, exposing Kafka over port 9092 to the host machine.

version: '3.8'

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.4.0
    hostname: zookeeper
    container_name: zookeeper
    ports:
      - "2181:2181"
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  broker:
    image: confluentinc/cp-server:7.4.0
    hostname: broker
    container_name: broker
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
      - "9101:9101"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_METRIC_REPORTERS: io.confluent.metrics.reporter.ConfluentMetricsReporter
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      CONFLUENT_METRICS_REPORTER_BOOTSTRAP_SERVERS: broker:29092
      CONFLUENT_SUPPORT_CUSTOMER_ID: 'anonymous'

Enter fullscreen mode Exit fullscreen mode

Docker CLI Commands to Spin Up Infrastructure

To start the background message broker infrastructure, navigate to your project directory containing the configuration file and run:

# Start Kafka and Zookeeper services in detached mode
docker compose up -d

# Verify that your containers are running normally
docker ps

Enter fullscreen mode Exit fullscreen mode

Data Ingestion: The Producer Layer (producer.py)

The producer.py script handles data ingestion. It loops continuously through an array of nine target international cities (Nairobi, Accra, Cape Town, Riga, Brussels, Moscow, Seoul, London, and Sucre), requests their real-time weather information via RapidAPI, and dispatches the payload to the weather_raw Kafka topic.

A standout feature here is the automatic inline JSON serialization using a lambda function passed directly into the KafkaProducer constructor.

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

load_dotenv()

API_KEY  = os.getenv("API_KEY")
API_HOST = os.getenv("API_HOST")
LANG     = os.getenv("LANG")

topic = "weather_raw"

# Initialize Kafka Producer with integrated JSON byte-serializer
producer = KafkaProducer(
    bootstrap_servers = 'localhost:9092',
    value_serializer = lambda v : json.dumps(v).encode('utf-8')
)

Cities = ['Nairobi','Accra','Cape Town','Riga','Brussels','Moscow','Seoul','London','Sucre']

while True :
    for city in Cities :
        url = f"https://{API_HOST}/city?city={city}&lang={LANG}"
        headers = {
            'x-rapidapi-host': API_HOST,
            'x-rapidapi-key' : API_KEY
        }

        try :
            response = requests.get(url, headers=headers)

            if response.status_code == 200 :
                weather_data = response.json()
                producer.send(topic, value=weather_data)
                print(f"Sent weather data for {city}")
            else :
                print(f"Failed for {city} : {response.status_code} ")

        except Exception as e :
            print(f"Error for {city} : {e}")

    producer.flush()
    time.sleep(1)

Enter fullscreen mode Exit fullscreen mode

To begin streaming live API payloads into your Kafka cluster, run the producer engine script in your terminal:

python3 producer.py

Enter fullscreen mode Exit fullscreen mode

Output:

Storage & Transformation Tier: The ETL Consumer (consumer.py)

The consumer layer implements a true ETL pattern. Instead of just printing raw bytes, it targets the weather_raw topic, decodes the stream, flattens the highly nested API structures into a uniform format using Pandas, and loads the records into a PostgreSQL database instance using SQLAlchemy.

from kafka import KafkaConsumer
from sqlalchemy import create_engine
from dotenv import load_dotenv
import os
import json
import pandas as pd

load_dotenv()

Postgres_URI = os.getenv("POSTGRES_URI")
engine = create_engine(Postgres_URI)

# Initialize Kafka Consumer with native byte-decoding
consumer = KafkaConsumer(
    'weather_raw',
    bootstrap_servers = 'localhost:9092',
    auto_offset_reset = 'earliest',
    enable_auto_commit = True,
    value_deserializer = lambda x : json.loads(x.decode('utf-8'))
)

print("Consumer started listening ...")

for message in consumer :
    try :
        data = message.value

        # EXTRACT & TRANSFORM: Defensive parsing handles nested JSON safely
        transformed_data = {
            "city"        : data.get("name"),
            "temperature" : data.get("main",{}).get("temp"),
            "humidity"    : data.get("main",{}).get("humidity"),
            "pressure"    : data.get("main",{}).get("pressure"),
            "weather"     : data.get("weather",[{}])[0].get("main"),
            "description" : data.get("weather",[{}])[0].get("description"),
            "wind_speed"  : data.get("wind",{}).get("speed")
        }

        # Structure as a Pandas DataFrame
        df = pd.DataFrame([transformed_data])
        print("\n Transformed weather data")

        # LOAD: Persist metrics into the PostgreSQL destination table
        df.to_sql("weather_kafka", con=engine, if_exists="append", index=False)
        print(f"Loaded weather data for {transformed_data['city']}")

    except Exception as e :
        print(f"Consumer error : {e}")

Enter fullscreen mode Exit fullscreen mode

Open a separate terminal shell pane and launch the engine to begin populating your relational database rows in real-time:

python3 consumer.py

Enter fullscreen mode Exit fullscreen mode

Output:

Pipeline Verification & Data Verification

To verify the integration, we can monitor the execution traces across the python workflows, and query the final warehouse target to confirm data persistence.

Dual-Terminal Execution Log Comparison

When running both backend applications synchronously, your live terminal workspace layout matches the active message flow:

Verifying Records in the PostgreSQL Database

Because the consumer script implements df.to_sql(..., if_exists="append"), every iteration builds out relational records in real-time. Opening a connection tool or terminal CLI to your PostgreSQL instance reveals the transformed schemas waiting for analytics:

select * from weather_kafka;

Enter fullscreen mode Exit fullscreen mode

Output log view:

Conclusion

This project successfully establishes a production-grade blueprint for real-time streaming data pipelines. By combining Docker container isolation with Kafka's decoupled storage guarantees, the system handles data ingestion loops safely without threatening the state of the loading layer. Using Python, Pandas, and SQLAlchemy turns nested API variations into structured relational records, resulting in an automated, robust data engine ready for downstream business intelligence dashboards.