Параллельная работа двух сервисов с очередью: multiproecssing vs threading?

Есть вебсокет канал, в него прилетают чанки аудиозаписи. Нужно их склеивать, преобразовывать в файл, делать различные операции и т.п. Результат нужно отправлять в гугл. Результат от гугла - отправлять обратно на фронт.

Думаю это всё разделить на два сервиса и делать параллельно. С потоками/процессами раньше не работал. Поэтому прошу подсказать.

Вижу примерно следующее.

  1. создать первый сервис, в котором будут производится различные операции с чанками, файлом. Результат ложить в очередь.

  2. создать второй сервис, в котором будут доставаться данные из очереди, будут отправляться в гугл и результат уже - на фронт.

Пример вебсокет консюмера:

class ConnectionManager:
    def __init__(self):
        self.active_connections: List[WebSocket] = []
        self.words_list: Dict[str, dict] = {}
        
    async def connect(self, websocket: WebSocket):
        await websocket.accept()
        if websocket not in self.active_connections:
            self.active_connections.append(websocket)
        if str(websocket.url) not in self.words_list:
            self.words_list[str(websocket.url)] = {"chunks": "", "text": ""}

    async def disconnect(self, websocket: WebSocket):
        await websocket.send_text(self.words_list[str(websocket.url)]["text"])
        self.active_connections.remove(websocket)

    async def broadcast(self, message: str, websocket: WebSocket):
        new_chunks = merge_base64(
            self.words_list[str(websocket.url)]["chunks"], message
        )
        text = recognition_voice_data(new_chunks)
        if not text or self.words_list[str(websocket.url)]["text"].endswith(
            text
        ):
            for connection in self.active_connections:
                if connection.url == websocket.url:
                    await connection.send_text(text)
            self.words_list[str(websocket.url)]["text"] = ""
            self.words_list[str(websocket.url)]["chunks"] = ""
        else:
            self.words_list[str(websocket.url)]["text"] += text
            self.words_list[str(websocket.url)]["chunks"] = new_chunks

Вопросы:

  1. Как всё это организовать? Что нужно, процессы, или потоки? Погуглив, я пришёл к выводу, что мне нужны потоки. Так лиэто?

  2. Каким образом второй сервис будет понимать, что в очереди появились новые данные?

В качество очереди думаю использовать питоновский queue


Ответы (0 шт):