Design for 100+ concurrent users on a Whisper Live Transcription API with prompt query engine.
| Bottleneck | Impact |
|---|---|
asyncio.Semaphore(2) |
Max 2 concurrent transcriptions |
In-memory dict[str, Session] |
Lost on restart, no persistence |
| Single-process FastAPI + uvicorn | No horizontal scaling |
| Whisper.cpp subprocess per request | Heavy CPU/GPU per call |
| No message queue | Requests dropped under load |
graph TB
subgraph Clients
WEB[Browser<br/>mqtt.js WSS]
MOB[Mobile App<br/>Native MQTT]
IOT[Device<br/>MQTT]
end
subgraph Edge
LB[Load Balancer<br/>TLS termination]
ES_ING[API Gateway<br/>3-5 FastAPI pods]
end
subgraph Messaging
MQ[EMQX Cluster<br/>3-5 nodes<br/>open-source]
KV[Kafka Bridge<br/>EMQX built-in plugin]
end
subgraph Stream
K[Kafka<br/>3 brokers]
K_TJ["transcription.jobs<br/>(50 partitions)"]
K_TR["transcription.results<br/>(50 partitions)"]
end
subgraph Compute
WP[Whisper Workers<br/>10-50 GPU pods]
QW[Query Workers<br/>prompt + window processing]
end
subgraph Storage
S3[(S3/Object Store<br/>raw audio)]
ES[(Elasticsearch<br/>transcript index<br/>6 shards)]
PG[(PostgreSQL<br/>session history)]
RD[(Redis<br/>hot session state)]
end
subgraph Batch
AF[Airflow<br/>re-transcribe DAGs]
end
WEB -- WSS --> LB
MOB -- WSS --> LB
IOT -- MQTT --> LB
LB --> MQ
LB --> ES_ING
ES_ING -- produce --> K_TJ
K_TJ --> WP
WP --> S3
WP -- produce --> K_TR
K_TR -- KV bridge --> MQ
K_TR -- index --> ES
K_TR -- insert --> PG
K_TR -- update --> RD
MQ --> WEB
MQ --> MOB
MQ --> IOT
WEB -- "query/prompt (REST)" --> ES_ING
ES_ING -- ES search --> ES
ES_ING --> RD
ES_ING --> PG
AF -.-> WP
AF -.-> ES
transcription/{session_id}/result ← Worker → Kafka → EMQX → Client
transcription/{session_id}/progress ← progress percent updates
transcription/{session_id}/match ← prompt match found via ES
query/{query_id}/result ← synchronous prompt query via REST
Clients subscribe to transcription/{session_id}/+ -- the + wildcard matches any subtopic (result, progress, match). EMQX handles fan-out to all connected clients automatically. No separate WebSocket mesh needed -- the MQTT broker IS the mesh.
Why Kafka over RabbitMQ:
| Requirement | Kafka | RabbitMQ |
|---|---|---|
| Partition-based parallelism | Native | Requires shovels |
| Replay/retention | Configurable (days/weeks) | Consumer-acked delete |
| Throughput | Millions/sec | Tens of thousands/sec |
| Ordered processing per session | Partition key = session_id | Complex |
Kafka Topics:
topics:
transcription.jobs:
partitions: 50 # = max parallelism
replication: 3
retention.ms: 604800000 # 7 days for replay
key: session_id # ensures ordering per session
transcription.results:
partitions: 50
compaction: true # keep latest per key
transcription.windows:
partitions: 20
cleanup: delete
prompt.queries:
partitions: 10
retention: 3600000 # 1 hourConsumer Groups:
whisper-workers-N → transcription.jobs (N = 10-50, one per GPU)
prompt-workers → prompt.queries
window-workers → transcription.windows
Worker Architecture:
graph TB
subgraph Worker_Pod
KC[Kafka Consumer<br/>jobs topic]
JQ[Job Queue<br/>max 4 buffered]
GS[GPU Scheduler]
W1[whisper instance 1]
W2[whisper instance 2]
WN[whisper instance N]
KP[Kafka Producer<br/>results topic]
KC --> JQ --> GS
GS --> W1
GS --> W2
GS --> WN
W1 --> KP
W2 --> KP
WN --> KP
end
Scaling Calculation:
| Model | VRAM | Time per 60s audio | Workers per GPU |
|---|---|---|---|
| tiny (~39M params) | ~1 GB | ~2s real-time | 8-12 concurrent |
| base (~74M) | ~1.5 GB | ~4s | 4-6 |
| small (~244M) | ~2.5 GB | ~8s | 2-3 |
| medium (~769M) | ~5 GB | ~20s | 1 |
For 100 concurrent users uploading 60s audio:
- tiny model: 100 * 2s = 200s work time -> ~10 workers (2 GPUs)
- small model: 100 * 8s = 800s work time -> ~40 workers (14 GPUs)
Redis (hot state, fast access):
session:{id} -> { created_at, transcription_count }
user:{id}:active_session -> session_id
PostgreSQL (cold storage, history):
CREATE TABLE sessions (
id UUID PRIMARY KEY,
user_id UUID,
created_at TIMESTAMPTZ DEFAULT NOW(),
status TEXT DEFAULT 'active'
);
CREATE TABLE transcripts (
id BIGSERIAL PRIMARY KEY,
session_id UUID REFERENCES sessions(id),
text TEXT,
language TEXT,
duration_ms FLOAT,
segments JSONB,
created_at TIMESTAMPTZ DEFAULT NOW(),
text_search TSVECTOR GENERATED ALWAYS AS (
to_tsvector('english', text)
) STORED
);
CREATE INDEX idx_transcripts_session ON transcripts(session_id);
CREATE INDEX idx_transcripts_search ON transcripts USING GIN(text_search);Why Elasticsearch over PostgreSQL FTS alone:
| Feature | Elasticsearch | PostgreSQL FTS |
|---|---|---|
| Real-time index updates | Near-real-time (~1s) | Lock on write |
| Phrase queries with slop | Native match_phrase |
Complex tsquery |
| Highlighting | Built-in | Manual |
| Relevance tuning | BM25 + boosts | TF-IDF only |
| Aggregation speed | Sub-second on 10M docs | Seconds |
| Scaling | Horizontal sharding | Read replicas |
Elasticsearch Index Mapping:
{
"index": "transcripts",
"settings": {
"number_of_shards": 6,
"number_of_replicas": 2,
"analysis": {
"analyzer": {
"transcript_analyzer": {
"type": "custom",
"tokenizer": "standard",
"filter": ["lowercase", "stop", "snowball"]
}
}
}
},
"mappings": {
"properties": {
"session_id": { "type": "keyword" },
"text": { "type": "text", "analyzer": "transcript_analyzer" },
"segments": {
"type": "nested",
"properties": {
"start_ms": { "type": "float" },
"end_ms": { "type": "float" },
"text": { "type": "text", "analyzer": "transcript_analyzer" }
}
},
"language": { "type": "keyword" },
"created_at": { "type": "date" }
}
}
}Query: match prompt with time window:
{
"query": {
"bool": {
"must": [
{ "match": { "segments.text": "revenue growth" } },
{ "range": { "segments.start_ms": { "gte": 10000 } } },
{ "range": { "segments.end_ms": { "lte": 60000 } } }
]
}
},
"highlight": {
"fields": { "segments.text": {} }
}
}graph LR
subgraph Clients
WEB[Browser<br/>mqtt.js WSS]
MOB[Mobile App<br/>Native MQTT]
IOT[Device<br/>MQTT]
end
subgraph EMQX_Cluster
EMQ1[EMQX Node 1]
EMQ2[EMQX Node 2]
EMQ3[EMQX Node 3]
end
subgraph Backend
KV[Kafka Bridge Plugin<br/>emqx_bridge_kafka]
K[Kafka Cluster]
end
WEB -- WSS --> EMQ1
WEB -- WSS --> EMQ2
MOB -- WSS --> EMQ3
IOT -- MQTT --> EMQ1
EMQ1 --> KV
EMQ2 --> KV
EMQ3 --> KV
KV --> K
K --> KV
KV --> EMQ1
KV --> EMQ2
KV --> EMQ3
Why MQTT over raw WebSocket:
| Feature | MQTT over WSS | Raw WebSocket Mesh |
|---|---|---|
| QoS levels | 0, 1, 2 (exactly-once) | Must implement yourself |
| Retained messages | Reconnect gets last result immediately | Lost on disconnect |
| Wildcard subscriptions | session/+/result |
Must list each topic |
| Wire overhead | 2-byte minimum header | Multi-byte framing |
| Backpressure | Built-in flow control | Must build yourself |
| Browser support | mqtt.js library (~50KB) | Native WebSocket API |
| Mobile/IoT | Excellent (native libs, low power) | Must build reconnection |
| Server complexity | EMQX cluster (self-managed) | Lightweight, no broker |
| Latency | ~1-2ms broker overhead | Sub-millisecond direct |
When MQTT is the right choice:
- Mobile clients (native MQTT libraries are production-grade)
- Need QoS 1/2 for reliable delivery
- Want retained messages so reconnecting clients auto-catch-up
- Have IoT or embedded devices
- Need built-in auth (JWT, TLS client certs) without coding
When raw WS mesh is simpler:
- Browser-only, <10K connections
- Want to keep stack minimal (Redis pub/sub + WS)
- Sub-millisecond latency requirement
- Small team, don't want to run a broker
EMQX is Apache 2.0 open source. Both editions:
| Edition | License | Key features |
|---|---|---|
| Community | Apache 2.0 (free) | Full MQTT v3.1/v5, WSS, Kafka bridge, 1M+ concurrent connections, Kubernetes Helm chart |
| Enterprise | Commercial | Hot-hot replication, blue-green deploy, data persistence |
Self-management is straightforward:
# Helm chart -- one command for a 3-node cluster
helm repo add emqx https://repos.emqx.io/charts
helm install emqx-cluster emqx/emqx --set replicaCount=3No ZooKeeper/etcd needed -- EMQX auto-discovers nodes via DNS. The Kafka bridge is a built-in plugin (emqx_bridge_kafka), not a separate connector. Typical sizing: 3 x c5.xlarge (4 vCPU, 8GB) handles 500K+ concurrent MQTT connections.
Downside: Erlang/OTP under the hood. If something breaks at the broker level, debugging requires Erlang knowledge. The Helm chart abstracts most operations, but it's worth knowing.
| Workload | Scheduler | Why |
|---|---|---|
| Real-time API + Workers | Kubernetes | HPA, health checks, rolling updates, spot instances |
| Batch jobs | Airflow | DAG dependencies, retries, SLA tracking |
Airflow DAGs:
# Re-transcribe with better model
with DAG("upgrade_transcription", schedule="@weekly"):
fetch_old_audio >> whisper_large >> compare_results >> index_es
# Batch window analysis over all sessions
with DAG("daily_window_analysis"):
query_all_sessions >> compute_windows >> store_results >> notifysequenceDiagram
participant C as Client (mqtt.js)
participant LB as Load Balancer
participant MQ as EMQX
participant API as API Gateway
participant K as Kafka
participant W as Worker
participant ES as Elasticsearch
Note over C,ES: Step 1: Upload audio for transcription
C->>LB: POST /api/transcribe/stream/prompt<br/>audio file + prompt text
LB->>API: Forward
API->>K: Produce transcription.jobs<br/>{session_id, audio_ref, prompt}
Note over C,ES: Step 2: Worker transcribes
K->>W: Consume job
W->>W: Download audio from S3
W->>W: Run whisper.cpp
W->>K: Produce transcription.results<br/>{text, segments, language}
Note over C,ES: Step 3: Kafka bridge fans out to EMQX
K-->>MQ: Kafka Bridge Plugin<br/>forwards to MQTT topics
MQ-->>C: transcription/{sid}/result<br/>{text, segments}
Note over C,ES: Step 4: Index for search
K-->>ES: Index segments
Note over C,ES: Step 5: Prompt match
ES-->>ES: Search: match prompt<br/>against segments
ES-->>MQ: matches found
MQ-->>C: transcription/{sid}/match<br/>{matches: [{text, score, start_ms, end_ms}]}
Note over C,ES: Step 6: Client renders
C->>C: Render match cards<br/>with clickable timestamps
graph TB
subgraph "Kubernetes Cluster (15-30 nodes)"
LB[Load Balancer<br/>Traefik/NLB]
subgraph "API Tier (HPA 3-5)"
API1[FastAPI Pod 1]
API2[FastAPI Pod 2]
API3[FastAPI Pod 3]
end
subgraph "EMQX Tier (3 nodes)"
EMQ1[EMQX Pod 1]
EMQ2[EMQX Pod 2]
EMQ3[EMQX Pod 3]
end
subgraph "GPU Node Pool (p3.2xlarge spot)"
W1[Worker Pod 1<br/>V100]
W2[Worker Pod 2<br/>V100]
WN[Worker Pod N<br/>V100]
end
subgraph "Data Tier"
K1[Kafka Broker 1]
K2[Kafka Broker 2]
K3[Kafka Broker 3]
ES1[ES Master Node]
ES2[ES Data Node 1]
ES3[ES Data Node 2]
R1[Redis Primary]
R2[Redis Replica]
PG1[PostgreSQL Primary]
PG2[PostgreSQL Replica]
end
end
subgraph "External"
IN[Internet]
S3[(S3 Bucket)]
end
IN --> LB
LB --> API1 & API2 & API3
LB --> EMQ1 & EMQ2 & EMQ3
API1 & API2 & API3 --> K1 & K2 & K3
EMQ1 & EMQ2 & EMQ3 --> K1 & K2 & K3
K1 & K2 & K3 --> W1 & W2 & WN
W1 & W2 & WN --> S3
W1 & W2 & WN --> K1 & K2 & K3
K1 & K2 & K3 --> ES1 --> ES2 & ES3
K1 & K2 & K3 --> PG1 --> PG2
K1 & K2 & K3 --> R1 --> R2
API1 & API2 & API3 --> ES1
API1 & API2 & API3 --> PG1
API1 & API2 & API3 --> R1
| Component | Winner | Runner-up | Decision rationale |
|---|---|---|---|
| Message Queue | Kafka | RabbitMQ | Partition parallelism, replay retention, throughput |
| Search Engine | Elasticsearch | PostgreSQL FTS | Real-time indexing, highlighting, BM25 relevance |
| Real-time Push | MQTT over WSS (EMQX) | Raw WebSocket Mesh | QoS, retained msgs, mobile native, built-in auth |
| State Cache | Redis Cluster | Memcached | Pub/Sub channels, data structures, persistence |
| Batch Orchestration | Airflow | Prefect | Mature Kubernetes executor, DAG-as-code |
| ASR Engine | whisper.cpp | faster-whisper | C++ perf, minimal VRAM, no Python GIL |
| Container Orchestration | Kubernetes | Nomad | HPA, GPU scheduling, ecosystem maturity |
| Scenario | Recommendation |
|---|---|
| Browser-only, <10K connections | Raw WS mesh (Socket.IO + Redis adapter) -- simpler, no broker |
| Mobile + Web + IoT | MQTT over WSS -- one protocol fits all |
| Need at-least-once delivery | MQTT (QoS 1) -- built-in, no custom ack logic |
| Reconnecting clients must catch up | MQTT retained messages -- auto-send last result |
| Sub-millisecond latency required | Raw WS mesh -- no broker hop |
| Small team, minimal ops | Raw WS mesh -- nothing extra to operate |
| Production at scale | MQTT over WSS -- EMQX handles auth, clustering, monitoring |
# values.yaml for EMQX Helm chart
replicaCount: 3
config:
listeners:
wss:
bind: "0.0.0.0:8084"
# TLS terminates at load balancer; WSS inside cluster
mqtt:
bind: "0.0.0.0:1883"
bridges:
kafka:
enable: true
servers: "kafka-0.kafka-headless:9092,kafka-1.kafka-headless:9092"
topics:
- topic: "transcription.results"
direction: "ingress" # Kafka -> EMQX
- topic: "transcription.jobs"
direction: "egress" # EMQX -> Kafka (for future use)import mqtt from 'mqtt'
const client = mqtt.connect('wss://transcribe.example.com/mqtt', {
clientId: `web-${sessionId}`,
clean: false, // persist session
})
client.subscribe(`transcription/${sessionId}/+`, { qos: 1 })
client.on('message', (topic, payload) => {
const event = JSON.parse(payload.toString())
const subtopic = topic.split('/').pop()
if (subtopic === 'result') {
renderTranscription(event.text, event.segments)
} else if (subtopic === 'match') {
renderMatches(event.matches)
} else if (subtopic === 'progress') {
updateProgress(event.percent)
}
})# Current run_whisper -> Kafka consumer handler
async def handle_transcription_job(msg):
audio = await download_from_s3(msg.audio_ref)
result = await run_whisper(audio, msg.language)
await kafka_producer.send("transcription.results", result)
await index_elasticsearch(result)# Before: sessions: dict[str, Session]
# After:
await redis.set(f"session:{sid}", session_json, ex=86400)
await pg.execute("INSERT INTO transcripts (...) VALUES (...)")# Before: search_prompt(segments, prompt) - Python TF-IDF
# After:
result = await es.search(index="transcripts", body={
"query": {"match": {"segments.text": prompt}}
})# Before: sessions[sid].active_websockets.append(ws)
# After:
# Workers produce to Kafka -> EMQX bridge auto-delivers
# No direct WebSocket management in API code
# Clients connect to EMQX via WSS and subscribe to topics| Component | Spec | Monthly |
|---|---|---|
| 2x GPU nodes (p3.2xlarge spot) | 1x V100 each | ~$640 |
| 3x API nodes (c5.xlarge) | 4 vCPU, 8GB | ~$400 |
| 3x EMQX (c5.xlarge) | 4 vCPU, 8GB | ~$400 |
| 3x Kafka (m5.large) | 2 vCPU, 8GB + 500GB EBS | ~$500 |
| 3x ES data (m5.xlarge) | 4 vCPU, 16GB + 500GB EBS | ~$600 |
| 1x PostgreSQL (db.r5.large) | 2 vCPU, 16GB | ~$200 |
| Redis (cache.r5.large x 2) | 13GB each | ~$300 |
| S3 + transfer (~5TB) | ~$150 | |
| Total (spot / on-demand) | ~$3,190 / ~$5,690 | |
| Per user/month | ~$32 / ~$57 |
Note: EMQX replaces the WebSocket mesh tier in the cost breakdown. The 3x EMQX nodes serve both as the real-time delivery layer and eliminate the need for a separate custom WS mesh service.