Skip to content

Latest commit

 

History

History
648 lines (517 loc) · 17.7 KB

File metadata and controls

648 lines (517 loc) · 17.7 KB

Distributed Transcription System Architecture

Design for 100+ concurrent users on a Whisper Live Transcription API with prompt query engine.


1. Current Bottlenecks

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

2. Complete Architecture (MQTT over WSS)

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
Loading

Topic Design

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.


3. Component Deep-Dive

3.1 Message Queue: Kafka

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 hour

Consumer Groups:

whisper-workers-N    → transcription.jobs   (N = 10-50, one per GPU)
prompt-workers       → prompt.queries
window-workers       → transcription.windows

3.2 Worker Pool: whisper.cpp + GPU

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
Loading

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)

3.3 State & Session: Redis + PostgreSQL

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);

3.4 Search: Elasticsearch

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": {} }
  }
}

3.5 Real-time Delivery: MQTT over WSS via EMQX

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
Loading

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 -- Open Source, Self-Managed

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=3

No 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.


3.6 Orchestration: K8s + Airflow

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 >> notify

4. Data Flow: Real-time Transcription + Prompt Match

sequenceDiagram
    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
Loading

5. Deployment Topology

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
Loading

6. Component Trade-offs

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

MQTT over WSS vs Raw WebSocket Mesh: When to Choose

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

7. EMQX Configuration

Helm Deployment

# 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)

Client Connection (Browser)

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)
  }
})

8. Migration Path

Phase 1: Extract + Buffer

# 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)

Phase 2: Persist State

# Before: sessions: dict[str, Session]
# After:
await redis.set(f"session:{sid}", session_json, ex=86400)
await pg.execute("INSERT INTO transcripts (...) VALUES (...)")

Phase 3: Scale Search

# Before: search_prompt(segments, prompt) - Python TF-IDF
# After:
result = await es.search(index="transcripts", body={
    "query": {"match": {"segments.text": prompt}}
})

Phase 4: Replace WebSocket with EMQX

# 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

9. Cost Estimate (100 users, 60s audio each)

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.