DEV Community

Cover image for Building a Real-Time Social Listening & Brand Sentiment Engine with Python, Gemini 1.5, and DuckDB
Ruesch Manny
Ruesch Manny

Posted on Originally published at mannyverse767.gumroad.com

Building a Real-Time Social Listening & Brand Sentiment Engine with Python, Gemini 1.5, and DuckDB

Enterprise social listening platforms like Brandwatch or Sprout charge upwards of $800 to $2,000/month per seat. Under the hood, many still rely on brittle dictionary lookup tables (VADER-style positive/negative lexical scoring) or generic regex filters. When a developer tweets "Oh fantastic, another breaking API change deployed straight to prod on a Friday!", classical sentiment classifiers mark it as overwhelmingly positive due to "fantastic".

Meanwhile, writing your own real-time monitoring service often degrades into a distributed systems mess: rate limits across Reddit and Hacker News, unstructured LLM hallucinations, deduplication bottlenecks, and ballooning database costs.

Here is how to build an asynchronous, production-ready brand intelligence engine using Python, Google's Gemini 1.5 Flash structured outputs, and DuckDB—running for less than $1/month in API costs.


1. The Architecture

The engine is designed around a decoupled, non-blocking ingestion and evaluation pattern:

graph TD
    A[Reddit API / HN Firebase / Bluesky Firehose / RSS] -->|Async Polling| B[Ingestion Layer: Asyncio + HTTPX]
    B -->|Content Hash SHA-256| C[(DuckDB Local Store - Deduplication)]
    C -->|New Items| D[Gemini 1.5 Flash Engine]
    D -->|Pydantic Structured JSON| E[Alert Dispatcher]
    E -->|Urgency >= High or Polarity <= -0.6| F[Slack / Discord Webhook]
    E -->|Daily Rollup & Trend Metrics| G[Parquet / Analytical Store]

Core Pillars:

  1. Async Multi-Source Ingestion: Pulls from Reddit (asyncpraw), Hacker News (Firebase REST API), Bluesky (AT Protocol firehose), and custom RSS feeds asynchronously.
  2. Deduplication Engine: Generates an identity fingerprint (sha256(platform + post_id)) checked against local embedded DuckDB before hitting LLM inference.
  3. Structured Gemini 1.5 Flash Parsing: Uses schema-constrained JSON outputs to extract polarity (-1.0 to 1.0), urgency classification (low, medium, critical), entity references, and sarcasm detection.
  4. Downstream Alert Dispatcher: Routes critical events directly to Slack/Discord webhooks while appending all metadata into DuckDB for fast OLAP rollups.

2. Ingestion & Deduplication via DuckDB

Rather than spinning up Postgres or Redis for transient social data, we use DuckDB. DuckDB runs in-process, provides zero-latency ACID writes, and allows vector-like analytical SQL on JSON columns.

import duckdb
import hashlib

class StorageManager:
    def __init__(self, db_path: str = "brand_intel.duckdb"):
        self.con = duckdb.connect(db_path)
        self._init_schema()

    def _init_schema(self):
        self.con.execute("""
            CREATE TABLE IF NOT EXISTS raw_mentions (
                id VARCHAR PRIMARY KEY,
                platform VARCHAR,
                author VARCHAR,
                content VARCHAR,
                url VARCHAR,
                created_at TIMESTAMP,
                processed BOOLEAN DEFAULT FALSE
            );
            CREATE TABLE IF NOT EXISTS sentiment_results (
                id VARCHAR PRIMARY KEY,
                polarity FLOAT,
                urgency VARCHAR,
                is_sarcastic BOOLEAN,
                entities JSON,
                summary VARCHAR,
                evaluated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
                FOREIGN KEY (id) REFERENCES raw_mentions(id)
            );
        """)

    def is_duplicate(self, platform: str, external_id: str) -> bool:
        uid = hashlib.sha256(f"{platform}:{external_id}".encode()).hexdigest()
        res = self.con.execute("SELECT 1 FROM raw_mentions WHERE id = ?", [uid]).fetchone()
        return res is not None

    def insert_mention(self, platform: str, external_id: str, author: str, content: str, url: str, created_at):
        uid = hashlib.sha256(f"{platform}:{external_id}".encode()).hexdigest()
        self.con.execute("""
            INSERT INTO raw_mentions (id, platform, author, content, url, created_at)
            VALUES (?, ?, ?, ?, ?, ?)
            ON CONFLICT (id) DO NOTHING
        """, [uid, platform, author, content, url, created_at])
        return uid
Enter fullscreen mode Exit fullscreen mode

3. Structured Gemini 1.5 Flash Evaluation

To eliminate hallucinations and avoid parsing regex out of freeform LLM completions, we leverage Google GenAI's native Pydantic schema validation. We use gemini-1.5-flash due to its sub-second latency and $0.075 per 1M tokens cost tier.

The Evaluation Schema

from pydantic import BaseModel, Field
from typing import List, Literal
import google.generativeai as genai
import os

genai.configure(api_key=os.environ["GEMINI_API_KEY"])

class SentimentEvaluation(BaseModel):
    polarity: float = Field(
        description="Polarity score strictly between -1.0 (extremely negative/hostile) and 1.0 (enthusiastic/positive)."
    )
    urgency: Literal["low", "medium", "high", "critical"] = Field(
        description="'critical' indicates data breaches, service downtime, or legal threats; 'high' indicates severe bugs or churn threats; 'medium' is regular feedback; 'low' is passive mention."
    )
    is_sarcastic: bool = Field(description="True if the statement implies the opposite of its literal words.")
    entities: List[str] = Field(description="List of specific products, competitor names, or people mentioned.")
    summary: str = Field(description="1-sentence technical digest of the user's issue or praise.")
Enter fullscreen mode Exit fullscreen mode

Inference Function with Enforced JSON Schema

import json

def evaluate_mention(brand_name: str, text: str) -> SentimentEvaluation:
    prompt = f"""
    Analyze the following social mention referencing the brand: '{brand_name}'.
    Contextualize sarcasm, tech nuances, and developer sentiment.

    Content:
    """{text}"""
    """

    model = genai.GenerativeModel("gemini-1.5-flash")

    response = model.generate_content(
        prompt,
        generation_config=genai.GenerationConfig(
            response_mime_type="application/json",
            response_schema=SentimentEvaluation,
            temperature=0.1,
        ),
    )

    return SentimentEvaluation.model_validate_json(response.text)
Enter fullscreen mode Exit fullscreen mode

4. Threshold Alerting & Webhook Routing

When a negative spike occurs (e.g., a critical authentication vulnerability drops on Hacker News), latency matters. The dispatcher inspects the structured output and fires immediate payloads to Slack, Discord, or an incident management endpoint like PagerDuty.

import httpx
import asyncio

SLACK_WEBHOOK_URL = os.environ.get("SLACK_WEBHOOK_URL")

async def dispatch_alert(eval_data: SentimentEvaluation, platform: str, url: str, content: str):
    # Alert trigger: high/critical urgency OR heavily negative sentiment
    if eval_data.urgency in ["high", "critical"] or eval_data.polarity <= -0.5:
        color = "#FF0000" if eval_data.urgency == "critical" else "#FFA500"

        payload = {
            "attachments": [
                {
                    "color": color,
                    "title": f"🚨 Brand Alert [{eval_data.urgency.upper()}] on {platform}",
                    "title_link": url,
                    "fields": [
                        {"title": "Summary", "value": eval_data.summary, "short": False},
                        {"title": "Polarity", "value": f"{eval_data.polarity:.2f}", "short": True},
                        {"title": "Sarcasm Detected", "value": str(eval_data.is_sarcastic), "short": True},
                        {"title": "Entities", "value": ", ".join(eval_data.entities) or "None", "short": True},
                    ],
                    "text": f"*Snippet:*
> {content[:280]}...",
                }
            ]
        }

        async with httpx.AsyncClient() as client:
            resp = await client.post(SLACK_WEBHOOK_URL, json=payload)
            resp.raise_for_status()
Enter fullscreen mode Exit fullscreen mode

5. Deployment, Concurrency & Rate Limits

Running high-frequency pollers against varied APIs presents two main issues:

  1. Rate Limit Exhaustion: Reddit restricts standard OAuth to 60 requests/minute. Gemini 1.5 Flash has a free-tier limit of 15 RPM and Tier-1 pay-as-you-go limit of 1,000 RPM.
  2. Event Loop Starvation: Blocking HTTP calls freeze background listeners.

The Concurrency Pool Pattern

Use an asyncio.Queue paired with worker pools to guarantee your workers respect rate bounds:

async def analysis_worker(queue: asyncio.Queue, db: StorageManager, brand_name: str):
    while True:
        item = await queue.get()
        uid, platform, author, content, url = item
        try:
            # Non-blocking run for Gemini evaluation
            eval_res = await asyncio.to_thread(evaluate_mention, brand_name, content)

            # Save analytics to DuckDB
            db.con.execute("""
                INSERT INTO sentiment_results (id, polarity, urgency, is_sarcastic, entities, summary)
                VALUES (?, ?, ?, ?, ?, ?)
            """, [uid, eval_res.polarity, eval_res.urgency, eval_res.is_sarcastic, 
                  json.dumps(eval_res.entities), eval_res.summary])

            # Route Alerts
            await dispatch_alert(eval_res, platform, url, content)

        except Exception as e:
            print(f"[ERROR] Failed to process {uid}: {e}")
        finally:
            queue.task_done()
Enter fullscreen mode Exit fullscreen mode

Analytical Daily Rollups via DuckDB

To view aggregated brand health over the past 24 hours without loading heavy Pandas workflows:

SELECT 
    date_trunc('hour', evaluated_at) as hour_bucket,
    AVG(polarity) as avg_sentiment,
    COUNT(CASE WHEN urgency IN ('high', 'critical') THEN 1 END) as critical_count,
    COUNT(*) as total_mentions
FROM sentiment_results
WHERE evaluated_at >= NOW() - INTERVAL '24 HOURS'
GROUP BY 1
ORDER BY 1 DESC;
Enter fullscreen mode Exit fullscreen mode

6. Conclusion & Turnkey Template

With under 200 lines of modern Python, you can replace a legacy multi-thousand-dollar monitoring subscription with a real-time system that catches sarcasm, tags entities, writes to low-footprint local storage, and pings your on-call channel in real time.

You can manually implement this architecture using the snippets above, wire your own OAuth credentials, and maintain the worker pipelines yourself.

If you prefer to deploy a battle-tested, production-ready codebase immediately—complete with pre-configured connectors for Reddit, Bluesky, Hacker News, RSS feeds, an automated CLI dashboard, and native Docker Compose files:

Have questions about tuning the schema for specific B2B niche contexts? Drop a comment below!

Top comments (0)