> ## Documentation Index
> Fetch the complete documentation index at: https://docs.pyannote.ai/llms.txt
> Use this file to discover all available pages before exploring further.

# Live diarized transcription: per-speaker transcription streams (Approach 2)

> Learn how to route each pyannoteAI speaker to a dedicated transcription stream and merge the results into one live transcript.

## Why use one transcription stream per speaker?

Live captions, meeting assistants, and call analytics need to answer one question: who said what?

In this tutorial, we will combine pyannoteAI live diarization with one transcription stream per speaker. Each stream is assigned a pyannoteAI speaker label, so all text returned by that stream can be attributed directly to that speaker. The client combines timestamped results instead of matching transcript events to speaker turns.

Sending each speaker's speech to a separate transcription session may improve accuracy compared with sending all speakers to one session.

This approach requires a known upper bound on the number of speakers. The client opens one transcription stream for each possible speaker, so transcription cost scales with that upper bound. It also buffers audio while waiting for speaker events, which adds transcription latency. [Approach 1: single STT stream](/tutorials/live-diarized-transcription-single-stt-stream) uses one transcription session for all speakers, but requires reconciliation logic. Evaluate both approaches with your audio, cost, and latency requirements.

<Note>
  Read [Combining real-time diarization and transcription](https://pyannote.ai/blog/combining-real-time-diarization-and-transcription) for the concepts behind both designs. This tutorial implements Approach 2.
</Note>

## Build a live speaker-attributed transcript

We will build a command-line program that listens to your microphone and prints a live transcript:

```text theme={null}
[SPEAKER_00] Hey, are you seeing the numbers from last week?
[SPEAKER_01] Yeah, I just pulled them up. Revenue is up about twelve percent.
[SPEAKER_00] Nice, that is ahead of the forecast.
```

This example sets `max_speakers` to 2. It supports up to two speakers and opens two AssemblyAI sessions before capture starts. Both sessions incur transcription costs for the full call duration, even if the conversation has only one speaker.

## Prerequisites

* A pyannoteAI API key from the [dashboard](https://dashboard.pyannote.ai)
* An AssemblyAI API key from the [AssemblyAI dashboard](https://www.assemblyai.com/dashboard)
* Python 3.10+
* [uv](https://docs.astral.sh/uv/getting-started/installation/)
* Microphone access

The full script declares its dependencies inline. `uv run` installs them in an isolated environment automatically.

Create a `.env` file:

```bash theme={null}
PYANNOTEAI_API_KEY=sk_xxx
ASSEMBLYAI_API_KEY=xxxxxxxx
```

## Architecture overview

The script captures each microphone frame once. It sends the original float32 frame to pyannoteAI and stores the frame in a short buffer. The router uses speaker events to send PCM16 speech to each active speaker's AssemblyAI stream and silence to the other streams.

Speaker events arrive after the corresponding audio. The delayed audio buffer holds audio for one second to give those events time to arrive before routing. If the router processes speech before its speaker-start event arrives, it can send silence instead of speech. This fixed delay does not guarantee that all events arrive in time; adjust it for your network and processing latency.

```mermaid theme={null}
%%{init: {'themeVariables': {'fontSize': '12px'}, 'flowchart': {'nodeSpacing': 20, 'rankSpacing': 24}}}%%
flowchart TB
    microphone[Microphone] --> queue[Audio queue]
    queue --> pyannote[pyannoteAI stream]
    queue --> buffer[Delayed audio buffer]
    pyannote -->|Speaker events| router[Speaker router]
    buffer --> router
    router -->|Speaker audio or silence| aai0[AssemblyAI stream 0]
    router -->|Speaker audio or silence| aai1[AssemblyAI stream 1]
    router -->|Speaker audio or silence| aaiN[AssemblyAI stream N-1]
    aai0 --> timeline[Shared transcript timeline]
    aai1 --> timeline
    aaiN --> timeline
    timeline --> terminal[Terminal]
```

The diagram shows selected streams for a general `N`; this example opens two. All `N` AssemblyAI streams receive the same amount of audio. Silence keeps each stream's clock moving when its speaker is inactive. Word timestamps from all streams therefore use one common timeline.

<Warning>
  Routing is not source separation. If multiple speakers talk into one microphone at the same time, each active stream receives the same mixed frame. Use separate input channels or a source-separation model when you must isolate overlapping voices.
</Warning>

## 1. Choose the speaker upper bound and open a pyannoteAI stream

Set `max_speakers` to the maximum number of speakers that can join the conversation. This value is `N`, the speaker-count upper bound. The script must know it before opening the transcription sessions. The router assigns each new pyannoteAI speaker label to one of the `N` slots. pyannoteAI live diarization tracks up to eight speakers, so use a value from 1 to 8:

```python theme={null}
max_speakers = 2  # N, the known speaker-count upper bound


def create_pyannote_stream() -> tuple[str, str]:
    response = requests.post(
        "https://api.pyannote.ai/v1/live",
        headers={"Authorization": f"Bearer {PYANNOTEAI_API_KEY}"},
    )
    response.raise_for_status()
    data = response.json()
    return data["id"], data["url"]
```

The response provides a stream ID and a single-use WebSocket URL. pyannoteAI accepts 16 kHz mono float32 PCM in 100 ms frames and returns stable speaker labels with start and end timestamps.

The script holds audio for one second before routing it:

```python theme={null}
ROUTING_DELAY = 1.0
```

Speaker events describe audio that the microphone already captured. This delay gives those events time to arrive before the matching buffered frame is sent to AssemblyAI. Increase it if network or processing latency causes clipped turn starts.

## 2. Open N AssemblyAI streams

Create one AssemblyAI WebSocket for each of the `N` speaker slots before audio capture starts:

```python theme={null}
params = urlencode({
    "sample_rate": SAMPLE_RATE,
    "encoding": "pcm_s16le",
    "speech_model": "universal-3-5-pro",
    "speaker_labels": "false",
    "min_turn_silence": 100,
    "max_turn_silence": 1000,
    "continuous_partials": "true",
})
assembly_url = f"wss://streaming.assemblyai.com/v3/ws?{params}"

aai_streams = [
    await stack.enter_async_context(
        websockets.connect(
            assembly_url,
            additional_headers={"Authorization": ASSEMBLYAI_API_KEY},
        )
    )
    for _ in range(max_speakers)
]
```

`universal-3-5-pro` uses turn detection. `min_turn_silence` starts an end-of-turn check after a short pause, and `max_turn_silence` forces the turn to end after one second of silence. AssemblyAI speaker labels stay disabled because pyannoteAI supplies the speaker identity.

<Warning>
  AssemblyAI bills each open streaming session by its duration, including periods that contain silence. Each session is billed independently, so `N` speaker slots produce `N` times the billed session duration of one stream. Always terminate every session when capture stops.
</Warning>

## 3. Route audio to each speaker

Store the pyannoteAI events with their timestamps:

```python theme={null}
def add_event(self, timestamp: float, speaker: str, started: bool) -> None:
    if started and speaker not in self.speaker_indexes:
        if len(self.speaker_indexes) < max_speakers:
            self.speaker_indexes[speaker] = len(self.speaker_indexes)
    self.events.append((timestamp, speaker, started))
    self.events.sort(key=lambda event: event[0])
```

For each buffered frame, allocate one row per speaker in a two-dimensional array. Copy samples to every row whose speaker is active. Leave the other rows at zero:

```python theme={null}
outputs = np.zeros(
    (max_speakers, len(frame.samples)),
    dtype=frame.samples.dtype,
)

for speaker in self.active:
    index = self.speaker_indexes.get(speaker)
    if index is not None:
        outputs[index, start_sample:end_sample] = frame.samples[start_sample:end_sample]
```

Convert each routed float32 array to little-endian PCM16 before sending it to AssemblyAI:

```python theme={null}
for stream, output in zip(aai_streams, outputs):
    pcm16 = (np.clip(output, -1, 1) * 32767).astype("<i2").tobytes()
    await stream.send(pcm16)
```

Each stream receives every frame. For example, while only `SPEAKER_00` talks, stream 0 receives the microphone samples and the other `N - 1` streams receive the same number of zero samples.

## 4. Handle partials, finals, and the shared timeline

AssemblyAI emits a `Turn` message more than once as a turn develops. Its `transcript` field contains the current text for the full turn, so each partial replaces the prior partial from that stream:

```python theme={null}
if not message["end_of_turn"]:
    timeline.update_partial(speaker, start, message["transcript"])
```

When `end_of_turn` is true, store the completed turn and remove that speaker's partial:

```python theme={null}
timeline.finalize(
    speaker,
    message["turn_order"],
    start,
    message["transcript"],
)
```

The first word's `start` value places the turn on the shared timeline. The terminal renderer sorts completed and partial turns by this timestamp. Results can arrive from the `N` AssemblyAI sessions in a different order, but their display order remains chronological.

## Full code

Save this complete script as `live_diarized_transcription_per_speaker.py`:

```python theme={null}
#!/usr/bin/env python3
# /// script
# requires-python = ">=3.10"
# dependencies = [
#     "sounddevice",
#     "numpy",
#     "requests",
#     "websockets>=13",
#     "python-dotenv",
# ]
# ///

import asyncio
import json
import os
import signal
from collections import deque
from contextlib import AsyncExitStack
from dataclasses import dataclass
from urllib.parse import urlencode

import numpy as np
import requests
import sounddevice as sd
import websockets
from dotenv import load_dotenv

load_dotenv()
PYANNOTEAI_API_KEY = os.environ["PYANNOTEAI_API_KEY"]
ASSEMBLYAI_API_KEY = os.environ["ASSEMBLYAI_API_KEY"]

SAMPLE_RATE = 16_000
CHUNK_SAMPLES = 1600
CHUNK_SECONDS = CHUNK_SAMPLES / SAMPLE_RATE
max_speakers = 2  # N, the known speaker-count upper bound
ROUTING_DELAY = 1.0


@dataclass
class AudioFrame:
    start: float
    samples: np.ndarray

    @property
    def end(self) -> float:
        return self.start + len(self.samples) / SAMPLE_RATE


class SpeakerRouter:
    def __init__(self) -> None:
        self.frames: deque[AudioFrame] = deque()
        self.events: list[tuple[float, str, bool]] = []
        self.active: set[str] = set()
        self.speaker_indexes: dict[str, int] = {}
        self.latest_time = 0.0
        self.warned_about_limit = False

    def add_frame(self, frame: AudioFrame) -> None:
        self.frames.append(frame)
        self.latest_time = frame.end

    def add_event(self, timestamp: float, speaker: str, started: bool) -> None:
        if started and speaker not in self.speaker_indexes:
            if len(self.speaker_indexes) < max_speakers:
                self.speaker_indexes[speaker] = len(self.speaker_indexes)
            elif not self.warned_about_limit:
                print(f"More than {max_speakers} speakers detected; extra speakers are muted.")
                self.warned_about_limit = True
        self.events.append((timestamp, speaker, started))
        self.events.sort(key=lambda event: event[0])

    def speaker_for_index(self, index: int) -> str:
        for speaker, speaker_index in self.speaker_indexes.items():
            if speaker_index == index:
                return speaker
        return f"SPEAKER_{index:02d}"

    def apply_event(self, speaker: str, started: bool) -> None:
        if started:
            self.active.add(speaker)
            return
        self.active.discard(speaker)

    def copy_active(self, outputs: np.ndarray, samples: np.ndarray, start: int, end: int) -> None:
        for speaker in self.active:
            index = self.speaker_indexes.get(speaker)
            if index is not None:
                outputs[index, start:end] = samples[start:end]

    def route(self, frame: AudioFrame) -> list[bytes]:
        outputs = np.zeros(
            (max_speakers, len(frame.samples)),
            dtype=frame.samples.dtype,
        )

        while self.events and self.events[0][0] <= frame.start:
            _, speaker, started = self.events.pop(0)
            self.apply_event(speaker, started)

        sample = 0
        while self.events and self.events[0][0] < frame.end:
            timestamp, speaker, started = self.events.pop(0)
            event_sample = round((timestamp - frame.start) * SAMPLE_RATE)
            event_sample = max(sample, min(len(frame.samples), event_sample))
            self.copy_active(outputs, frame.samples, sample, event_sample)
            self.apply_event(speaker, started)
            sample = event_sample

        self.copy_active(outputs, frame.samples, sample, len(frame.samples))
        return [
            (np.clip(output, -1, 1) * 32767).astype("<i2").tobytes()
            for output in outputs
        ]

    async def flush(self, streams, force: bool = False) -> None:
        cutoff = self.latest_time - ROUTING_DELAY
        while self.frames and (force or self.frames[0].end <= cutoff):
            frame = self.frames.popleft()
            for stream, pcm16 in zip(streams, self.route(frame)):
                await stream.send(pcm16)
            if force:
                await asyncio.sleep(CHUNK_SECONDS)


class TranscriptTimeline:
    def __init__(self) -> None:
        self.finals: dict[tuple[str, int], tuple[float, str, str]] = {}
        self.partials: dict[str, tuple[float, str, str]] = {}

    def update_partial(self, speaker: str, start: float, text: str) -> None:
        self.partials[speaker] = (start, speaker, text)
        self.render()

    def finalize(self, speaker: str, order: int, start: float, text: str) -> None:
        self.partials.pop(speaker, None)
        self.finals[(speaker, order)] = (start, speaker, text)
        self.render()

    def render(self) -> None:
        rows = [*self.finals.values(), *self.partials.values()]
        rows.sort(key=lambda row: row[0])
        print("\033[H\033[J", end="")
        print("Listening - speak now (Ctrl+C to stop)\n")
        for row in rows:
            suffix = " ..." if row in self.partials.values() else ""
            print(f"[{row[1]}] {row[2]}{suffix}")


def create_pyannote_stream() -> tuple[str, str]:
    response = requests.post(
        "https://api.pyannote.ai/v1/live",
        headers={"Authorization": f"Bearer {PYANNOTEAI_API_KEY}"},
    )
    response.raise_for_status()
    data = response.json()
    return data["id"], data["url"]


def microphone(loop, queue: "asyncio.Queue[np.ndarray]") -> sd.InputStream:
    def enqueue(frame: np.ndarray) -> None:
        try:
            queue.put_nowait(frame)
        except asyncio.QueueFull:
            pass

    def callback(indata, frames, time, status):
        loop.call_soon_threadsafe(enqueue, indata[:, 0].copy())

    stream = sd.InputStream(
        samplerate=SAMPLE_RATE,
        channels=1,
        dtype="float32",
        blocksize=CHUNK_SAMPLES,
        callback=callback,
    )
    stream.start()
    return stream


async def receive_pyannote(
    stream,
    router: SpeakerRouter,
    stop: asyncio.Event,
    done: asyncio.Event,
) -> None:
    try:
        async for raw in stream:
            message = json.loads(raw)
            message_type = message.get("type")
            data = message.get("data", {})
            if message_type == "diarization_speaker_start":
                router.add_event(data["timestamp"], data["speaker"], True)
            if message_type == "diarization_speaker_end":
                router.add_event(data["timestamp"], data["speaker"], False)
            if message_type == "error":
                print(f"pyannoteAI error: {message.get('message')}")
                stop.set()
    finally:
        done.set()
        if not stop.is_set():
            stop.set()


def turn_start(message: dict, fallback: float) -> float:
    words = message.get("words", [])
    if not words:
        return fallback
    return words[0]["start"] / 1000


async def receive_assemblyai(
    stream,
    index: int,
    router: SpeakerRouter,
    timeline: TranscriptTimeline,
    stop: asyncio.Event,
    done: asyncio.Event,
) -> None:
    try:
        async for raw in stream:
            message = json.loads(raw)
            if message.get("type") == "Turn" and message.get("transcript", "").strip():
                speaker = router.speaker_for_index(index)
                text = message["transcript"].strip()
                fallback = max(0.0, router.latest_time - ROUTING_DELAY)
                start = turn_start(message, fallback)
                if message["end_of_turn"]:
                    timeline.finalize(speaker, message["turn_order"], start, text)
                else:
                    timeline.update_partial(speaker, start, text)
            if message.get("type") == "Termination":
                done.set()
            if message.get("type") == "Error":
                print(f"AssemblyAI error for stream {index}: {message}")
                stop.set()
    finally:
        done.set()
        if not stop.is_set():
            stop.set()


async def pump_audio(
    queue: "asyncio.Queue[np.ndarray]",
    pyannote_stream,
    assembly_streams,
    router: SpeakerRouter,
    stop: asyncio.Event,
) -> None:
    samples_sent = 0
    while not stop.is_set():
        try:
            samples = await asyncio.wait_for(queue.get(), timeout=0.1)
        except asyncio.TimeoutError:
            continue
        await pyannote_stream.send(samples.astype("<f4").tobytes())
        router.add_frame(AudioFrame(samples_sent / SAMPLE_RATE, samples))
        samples_sent += len(samples)
        await router.flush(assembly_streams)


async def main() -> None:
    loop = asyncio.get_running_loop()
    stop = asyncio.Event()
    loop.add_signal_handler(signal.SIGINT, stop.set)

    stream_id, pyannote_url = create_pyannote_stream()
    print(f"pyannoteAI stream: {stream_id}")

    params = urlencode({
        "sample_rate": SAMPLE_RATE,
        "encoding": "pcm_s16le",
        "speech_model": "universal-3-5-pro",
        "speaker_labels": "false",
        "min_turn_silence": 100,
        "max_turn_silence": 1000,
        "continuous_partials": "true",
    })
    assembly_url = f"wss://streaming.assemblyai.com/v3/ws?{params}"

    async with AsyncExitStack() as stack:
        pyannote_stream = await stack.enter_async_context(websockets.connect(pyannote_url))
        assembly_streams = [
            await stack.enter_async_context(
                websockets.connect(
                    assembly_url,
                    additional_headers={"Authorization": ASSEMBLYAI_API_KEY},
                )
            )
            for _ in range(max_speakers)
        ]

        queue: "asyncio.Queue[np.ndarray]" = asyncio.Queue(maxsize=50)
        router = SpeakerRouter()
        timeline = TranscriptTimeline()
        mic = microphone(loop, queue)
        timeline.render()

        pyannote_done = asyncio.Event()
        assembly_done = [asyncio.Event() for _ in assembly_streams]
        pyannote_task = asyncio.create_task(
            receive_pyannote(pyannote_stream, router, stop, pyannote_done)
        )
        assembly_tasks = [
            asyncio.create_task(
                receive_assemblyai(
                    stream,
                    index,
                    router,
                    timeline,
                    stop,
                    assembly_done[index],
                )
            )
            for index, stream in enumerate(assembly_streams)
        ]
        pump_task = asyncio.create_task(
            pump_audio(queue, pyannote_stream, assembly_streams, router, stop)
        )
        tasks = [pyannote_task, *assembly_tasks, pump_task]

        await stop.wait()
        mic.stop()
        mic.close()
        await pump_task
        await pyannote_stream.send(json.dumps({"type": "end_of_stream"}))
        await pyannote_done.wait()
        await router.flush(assembly_streams, force=True)
        for stream in assembly_streams:
            await stream.send(json.dumps({"type": "Terminate"}))
        await asyncio.gather(*(done.wait() for done in assembly_done))
        for task in tasks:
            task.cancel()
        await asyncio.gather(*tasks, return_exceptions=True)


if __name__ == "__main__":
    asyncio.run(main())
```

## Run the script

Run the file from the directory that contains `.env`:

```bash theme={null}
uv run live_diarized_transcription_per_speaker.py
```

Speak into your microphone. Partial text updates in place. Completed AssemblyAI turns remain on screen and sort by their first word timestamp.

Press `Ctrl+C` to stop. The script stops microphone capture, waits for pyannoteAI to finalize its speaker events, flushes the delayed audio, and waits for all `N` AssemblyAI sessions to terminate with their final turns.

## Next steps

* To compare both designs, read [Combining real-time diarization and transcription](https://pyannote.ai/blog/combining-real-time-diarization-and-transcription).
* To use one transcription stream with reconciliation logic and avoid paying for one session per speaker slot, see [Approach 1: single STT stream](/tutorials/live-diarized-transcription-single-stt-stream).
* Set `max_speakers` to the smallest reliable upper bound for your conversation. The script opens and pays for that many AssemblyAI sessions for the full call, including silent sessions. If more speakers join, the router mutes the extra speakers.
* pyannoteAI live diarization supports up to eight speakers. If more than eight speakers join, it may merge speakers under the same labels.
* For production use, measure pyannoteAI event latency on your network and set `ROUTING_DELAY` above a high percentile of that measurement.

You can use another streaming transcription provider if it accepts continuous audio or explicit silence for each speaker stream. Change the audio conversion, connection setup, and event handler to match that provider. Partials may append text or replace the complete partial, turn-final signals use provider-specific fields, and timestamps may use a session clock or start at the first speech frame. Keep every stream on one common audio clock before merging its results.
