Why real-time sentiment matters for brands
When a product launch triggers a surge of tweets, the first hour can decide whether a campaign is a hit or a miss. A 2023 study by Brandwatch reported that 62% of consumers form an opinion within the first 30 minutes of a brand’s social buzz. Missing that window means losing actionable insight, which is why engineers are turning to streaming pipelines that can score sentiment the instant a post appears.
Architecture overview
The solution stitches three proven components: Apache Kafka as the backbone for ingesting and buffering social media events, the OpenAI API to transform raw text into a sentiment score, and Grafana to turn those scores into live dashboards. Data flows from a Kafka producer that pulls tweets via the Twitter API, through a consumer that calls OpenAI’s chat/completions endpoint, and finally into a second topic that Grafana reads via a simple Prometheus exporter.
System requirements
All components run on Linux (Ubuntu 22.04 LTS is the reference). You need at least 4 CPU cores, 8 GB RAM, and 50 GB of SSD space for Kafka logs. OpenAI API access requires a valid API key and a billing arrangement that can handle roughly 10,000 requests per day for a medium‑size brand.
Installing Apache Kafka
First, download the latest binary, extract it, and start the broker. The following bash snippet shows the exact steps.
wget https://downloads.apache.org/kafka/3.5.1/kafka_2.13-3.5.1.tgz && \
tar -xzf kafka_2.13-3.5.1.tgz && \
cd kafka_2.13-3.5.1 && \
bin/zookeeper-server-start.sh config/zookeeper.properties & \
bin/kafka-server-start.sh config/server.propertiesVerify the installation with bin/kafka-topics.sh --list --bootstrap-server localhost:9092. You should see the default internal topics.
Producing social media events
Python’s tweepy library streams live tweets that match a hashtag. Each tweet is serialized as JSON and pushed to a Kafka topic called raw_tweets.
import json, os, tweepy, kafka
api_key = os.getenv("TWITTER_API_KEY")
api_secret = os.getenv("TWITTER_API_SECRET")
auth = tweepy.OAuth2BearerHandler(os.getenv("TWITTER_BEARER_TOKEN"))
client = tweepy.Client(bearer_token=os.getenv("TWITTER_BEARER_TOKEN"))
producer = kafka.KafkaProducer(bootstrap_servers=['localhost:9092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8'))
def on_tweet(tweet):
payload = {"id": tweet.id, "text": tweet.text, "created_at": str(tweet.created_at)}
producer.send('raw_tweets', value=payload)
for tweet in tweepy.Paginator(client.search_recent_tweets,
query="#YourBrand", tweet_fields=['created_at'],
max_results=100).flatten(limit=1000):
on_tweet(tweet)
Run the script in a screen session; it will keep feeding Kafka as long as the hashtag is active.
Consuming tweets and calling OpenAI
The consumer reads from raw_tweets, extracts the text field, and asks OpenAI for a sentiment classification. The model gpt-4o-mini returns a JSON object with sentiment (positive, neutral, negative) and a confidence score.
import os, json, kafka, openai
openai.api_key = os.getenv("OPENAI_API_KEY")
consumer = kafka.KafkaConsumer('raw_tweets',
bootstrap_servers=['localhost:9092'],
value_deserializer=lambda m: json.loads(m.decode('utf-8')),
auto_offset_reset='earliest',
enable_auto_commit=True)
producer = kafka.KafkaProducer(bootstrap_servers=['localhost:9092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8'))
def analyze(text):
response = openai.ChatCompletion.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": f"Classify the sentiment of this tweet and return JSON with fields sentiment and confidence: {text}"}],
temperature=0
)
return json.loads(response.choices[0].message.content)
for msg in consumer:
tweet = msg.value
result = analyze(tweet['text'])
enriched = {**tweet, **result}
producer.send('sentiment_tweets', value=enriched)
The loop processes roughly 200 tweets per minute on a modest VM, staying well within OpenAI’s rate limits for a paid tier.
Exporting metrics for Grafana
Grafana can read time‑series data directly from Prometheus, so we expose a tiny exporter that converts each enriched tweet into a gauge metric. The exporter runs on port 9100 and increments counters for each sentiment category.
from prometheus_client import start_http_server, Counter
import json, kafka
POSITIVE = Counter('tweets_positive_total', 'Number of positive tweets')
NEUTRAL = Counter('tweets_neutral_total', 'Number of neutral tweets')
NEGATIVE = Counter('tweets_negative_total', 'Number of negative tweets')
consumer = kafka.KafkaConsumer('sentiment_tweets',
bootstrap_servers=['localhost:9092'],
value_deserializer=lambda m: json.loads(m.decode('utf-8')))
start_http_server(9100)
for msg in consumer:
sentiment = msg.value.get('sentiment', 'neutral').lower()
if sentiment == 'positive':
POSITIVE.inc()
elif sentiment == 'negative':
NEGATIVE.inc()
else:
NEUTRAL.inc()
After starting the exporter, add a Prometheus data source in Grafana and create a dashboard with three single‑stat panels that show the live counts. Use the query rate(tweets_positive_total[1m]) for a per‑minute trend.
Deployment tips and scaling
For production, containerize each component with Docker and orchestrate with Kubernetes. A typical pod spec allocates 0.5 CPU and 512 MiB RAM for the consumer, while the Kafka broker benefits from a dedicated StatefulSet with persistent volume claims. Enable Kafka’s log compaction on the sentiment_tweets topic to keep only the latest sentiment per tweet ID, which reduces storage overhead by up to 70%.
Conclusion
By coupling Apache Kafka’s fault‑tolerant streaming with OpenAI’s language model and Grafana’s visual power, you can turn a chaotic firehose of social posts into a clear, actionable sentiment dashboard. The pipeline runs in real time, scales horizontally, and delivers business‑critical insight within seconds of a post appearing. Start small—track a single hashtag—and expand to multiple brands once you validate latency and cost.
Sources
Apache Kafka Documentation, OpenAI API Reference, Grafana Labs Tutorials
Author: Mahmut Sarıkaya — sarikayadev.com