Automatisation de la synthèse vocale CosyVoice : Script Python pour le traitement par lots haute performance

Architecture du Pipeline de Synthèse Audio

Générer manuellement des pistes vocales via l'API CosyVoice devient rapidement non scalable dès que le volume dépasse la dizaine de segments. Une solution robuste consiste à découper le workflow en trois couches : ingestion de données, exécution réseau concurrente et persistance structurée. Ce modèle permet de remplacer les interactions GUI répétitives par un moteur programmable, capable de gérer les échecs transitoires et de tracer chaque opération.

1. Environnement et Dépendacnes

L'infrastructure minimale requiert Python 3.9+, une connexion stable vers l'endpoint TTS, et les bibliothèques suivantes :

pip install requests pandas concurrent.futures

Ces outils couvrent le parsing de feuilles de calcul, les requêtes HTTP synchrones et l'ordonnancement threadé. Aucun framework lourd n'est nécessaire pour cette charge de travail orientée I/O.

2. Modélisation et Ingestion des Données

Il est préférable de formaliser les jobs avec des structures typées plutôt que des dictionnaires bruts. Le chargement se fait depuis un CSV séparateur ; chaque ligne est transformée en instance objet contenant les métadonnées de synthèse.

from dataclasses import dataclass
import pandas as pd

@dataclass(slots=True)
class AudioJob:
    job_id: str
    script: str
    voice_profile: str
    tempo_ratio: float
    dest_filename: str

def fetch_jobs(dataset_path: str) -> list[AudioJob]:
    raw_data = pd.read_csv(dataset_path, dtype={"id": str})
    jobs = []
    for _, row in raw_data.iterrows():
        jobs.append(AudioJob(
            job_id=row["id"],
            script=str(row.get("text_content", "")),
            voice_profile=row.get("speaker", "zhitian"),
            tempo_ratio=float(row.get("speed", 1.0)),
            dest_filename=row.get("filename", f"clip_{row['id']}.wav")
        ))
    print(f"[INFO] {len(jobs)} segments prêts à être traités.")
    return jobs

3. Orchestration de l'Appel REST

Chaque job doit être adressé individuellement vers le service CosyVoice. La fonction ci-dessous encapsule la sérialisation JSON, l'envoi POST et la séquence binaire de sauvegarde. Les erreurs HTTP sont converties en exceptions personnalisées pour une propagation uniforme dans le pipeline parallèle.

import os
import json
import requests
from typing import Tuple

class SynthesisFailure(Exception):
    pass

OUTPUT_ROOT = "./generated_clips"

def dispatch_synthesis(job: AudioJob, server_endpoint: str) -> str:
    os.makedirs(OUTPUT_ROOT, exist_ok=True)
    
    payload = {
        "text": job.script,
        "speaker": job.voice_profile,
        "rate": job.tempo_ratio,
        "format": "wav"
    }
    
    resp = requests.post(
        server_endpoint,
        json=payload,
        headers={"Content-Type": "application/json"},
        timeout=30
    )
    
    if resp.status_code != 200:
        raise SynthesisFailure(f"Erreur serveur ({resp.status_code}): {resp.text[:100]}")
        
    save_path = os.path.join(OUTPUT_ROOT, job.dest_filename)
    with open(save_path, "wb") as stream:
        stream.write(resp.content)
        
    return save_path

4. Exécution Concurrente et Collecte

Plutôt que d'enchaîner les appels linéairement, on exploite un thread pool pour overlap le temps d'attente réseau. Le schéma utilise as\_completed pour récupérer les résultats au fur et à mesure et maintenir un compteur de progression.

from concurrent.futures import ThreadPoolExecutor, as_completed

def run_parallel_batch(jobs: list[AudioJob], endpoint: str, pool_size: int = 8) -> dict:
    completed = []
    errors = []
    
    with ThreadPoolExecutor(max_workers=pool_size) as executor:
        future_map = {executor.submit(dispatch_synthesis, j, endpoint): j for j in jobs}
        
        for done_future in as_completed(future_map):
            target_job = future_map[done_future]
            try:
                path = done_future.result()
                completed.append((target_job.job_id, path))
            except Exception as exc:
                errors.append((target_job.job_id, str(exc)))
                
    return {"success": completed, "failures": errors}

5. Reprise Transparente et Journalisation

Les interruptions réseau ou les throttles temporaires nécessitent un mécanisme de retour arrière. Au lieu de dépendre de bibliothèques externes, on implémente un décorateur de redémarrage exponentiel. Couplé au module logging standard, il assure une traçabilité persistante.

import logging
import time
from functools import wraps

logging.basicConfig(
    filename="pipeline_trace.log",
    filemode="w",
    level=logging.INFO,
    format="%(asctime)s | %(levelname)s | %(message)s"
)
logger = logging.getLogger("tts_engine")

def resilient_wrapper(func, max_retries: int = 3, base_delay: float = 2.0):
    @wraps(func)
    def wrapper(*args, **kwargs):
        attempt = 0
        while attempt < max_retries:
            try:
                return func(*args, **kwargs)
            except (requests.exceptions.ConnectionError, requests.exceptions.Timeout):
                attempt += 1
                delay = base_delay * (2 ** (attempt - 1))
                logger.warning(f"Rupture connecteur | Job {args[0].job_id} | Retard {delay:.1f}s avant nouvelle tentative")
                time.sleep(delay)
        raise SynthesisFailure(f"Capacité maximale de récupération épuisée pour {args[0].job_id}")
    return wrapper

# Application du décorateur à la fonction de dispatch
dispatch_with_fallback = resilient_wrapper(dispatch_synthesis)

6. Point d'Entrée et Argumentation CLI

L'agrégation finale utilise argparse pour exposer une interface reproductible. Le script peut s'exécuter directement en ligne de commande sans modification du code source.

import argparse

def build_cli():
    parser = argparse.ArgumentParser(description="Moteur Batch CosyVoice")
    parser.add_argument("--sheet", required=True, help="Chemin du CSV source")
    parser.add_argument("--rpc", required=True, help="URL de l'endpoint TTS")
    parser.add_argument("--threads", type=int, default=6, help="Nombre deWorkers simultanés")
    parser.add_argument("--recovery", action="store_true", help="Activer la reprise exponentielle")
    return parser.parse_args()

if __name__ == "__main__":
    opts = build_cli()
    job_queue = fetch_jobs(opts.sheet)
    
    handler = dispatch_with_fallback if opts.recovery else dispatch_synthesis
    
    stats = run_parallel_batch(job_queue, opts.rpc, opts.threads)
    
    logger.info(f"Pipeline terminé | Réussites: {len(stats['success'])} | Échecs: {len(stats['failures'])}")
    print(f"\nExécution clôturée. Vérifiez ./pipeline_trace.log pour les détails.")

Pistes d'Extension Technique

  • Orchestrateur Distribué : Pour des volumes supérieurs à 50k segments, migrer vers Celery ou Airflow permet de découpler la production des consommations et de gérer les files de messages persistantes.
  • Validation Post-Generation : Intégrer une étape de vérification automatique (contrôle du header RIFF WAV, test de lecture FFmpeg, ou rétro-ingestion ASR) garantit l'intégrité acoustique avant mise en production.
  • Limitation de Débit : Implémenter des sémaphores ou des algorithmes Token Bucket empêche de saturer les quotas d'API et évite les blocages IP temporaires.
  • Métadonnées Dynamiques : Remplacer le CSV statique par un connector SQLAlchemy ou MongoDB permet d'alimenter le pipeline directement depuis un CMS ou un ERP existant.

Étiquettes: cosyvoice Python Pandas threading api-rest

Publié le 6 septembre à 09h25