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:
-
Async Multi-Source Ingestion: Pulls from Reddit (
asyncpraw), Hacker News (Firebase REST API), Bluesky (AT Protocol firehose), and custom RSS feeds asynchronously. -
Deduplication Engine: Generates an identity fingerprint (
sha256(platform + post_id)) checked against local embedded DuckDB before hitting LLM inference. -
Structured Gemini 1.5 Flash Parsing: Uses schema-constrained JSON outputs to extract polarity (
-1.0to1.0), urgency classification (low,medium,critical), entity references, and sarcasm detection. - 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
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.")
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)
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()
5. Deployment, Concurrency & Rate Limits
Running high-frequency pollers against varied APIs presents two main issues:
- 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.
- 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()
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;
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:
- Instant Access on Whop
-
Direct Download on Gumroad (Use promo code
EARLYBIRDfor 20% off)
Have questions about tuning the schema for specific B2B niche contexts? Drop a comment below!
Top comments (0)