#!/usr/bin/env python3 import argparse import multiprocessing as mp import os import sqlite3 import time from pathlib import Path DB = Path("data/recording_index.sqlite3") OUT = Path("data/transcripts") MODEL = os.environ.get("WHISPER_MODEL", "small") OUT.mkdir(parents=True, exist_ok=True) def worker(worker_id, jobs): from faster_whisper import WhisperModel print( f"[Worker {worker_id}] " f"starte Whisper {MODEL}", flush=True, ) model = WhisperModel( MODEL, device="cuda", compute_type="int8", num_workers=1, ) con = sqlite3.connect(DB) for job_id, audio_path, output_path in jobs: try: # Andere Worker können denselben Job nicht übernehmen, # solange wir ihn atomar auf RUNNING setzen. cur = con.execute(""" UPDATE recordings SET status = 'TRANSCRIBING' WHERE id = ? AND status = 'NEW' """, (job_id,)) con.commit() if cur.rowcount != 1: continue print( f"[W{worker_id}] " f"{audio_path}", flush=True, ) segments, info = model.transcribe( audio_path, language="de", beam_size=5, vad_filter=True, vad_parameters={ "min_silence_duration_ms": 500, }, ) text_parts = [] for segment in segments: text = segment.text.strip() if text: text_parts.append(text) text = "\n".join(text_parts).strip() output = Path(output_path) output.parent.mkdir( parents=True, exist_ok=True, ) output.write_text( text, encoding="utf-8", ) status = ( "TRANSCRIBED" if text else "NO_SPEECH" ) con.execute(""" UPDATE recordings SET status = ?, transcript_path = ? WHERE id = ? """, ( status, str(output), job_id, )) con.commit() except Exception as exc: print( f"[W{worker_id}] FEHLER " f"{audio_path}: {exc}", flush=True, ) con.execute(""" UPDATE recordings SET status = 'TRANSCRIPTION_ERROR' WHERE id = ? """, (job_id,)) con.commit() con.close() print( f"[Worker {worker_id}] fertig", flush=True, ) def main(): parser = argparse.ArgumentParser() parser.add_argument( "--workers", type=int, default=2, ) args = parser.parse_args() if args.workers < 1: raise SystemExit( "--workers muss >= 1 sein" ) con = sqlite3.connect(DB) con.row_factory = sqlite3.Row # Persönliche/interne Aufzeichnungen ausschließen. con.execute(""" UPDATE recordings SET status = 'EXCLUDED_INTERNAL' WHERE extension = '42' AND status IN ('NEW', 'TRANSCRIBING') """) # Sehr kurze Aufnahmen ausschließen. con.execute(""" UPDATE recordings SET status = 'TOO_SHORT' WHERE duration IS NOT NULL AND duration < 2 AND status IN ('NEW', 'TRANSCRIBING') """) con.commit() rows = con.execute(""" SELECT id, path, recording_id FROM recordings WHERE status = 'NEW' ORDER BY recorded_at, id """).fetchall() con.close() print("=" * 72) print("PARALLELER WHISPER-LAUF") print("=" * 72) print(f"Modell: {MODEL}") print(f"Worker: {args.workers}") print(f"Kandidaten: {len(rows)}") print("Device: CUDA") print("Compute: int8") print() if not rows: print("Keine neuen Aufnahmen.") return # Jobs deterministisch auf Worker verteilen. jobs = [[] for _ in range(args.workers)] for index, row in enumerate(rows): output = ( OUT / f"{row['id']}_{row['recording_id']}.txt" ) jobs[index % args.workers].append( ( row["id"], row["path"], str(output), ) ) processes = [] started = time.time() for worker_id, worker_jobs in enumerate( jobs, 1, ): if not worker_jobs: continue process = mp.Process( target=worker, args=( worker_id, worker_jobs, ), ) process.start() processes.append(process) try: for process in processes: process.join() except KeyboardInterrupt: print("\nAbbruch angefordert – Worker werden beendet ...", flush=True) for process in processes: if process.is_alive(): process.terminate() for process in processes: process.join(timeout=10) for process in processes: if process.is_alive(): process.kill() # Jobs, die beim Abbruch noch liefen, wieder freigeben. cleanup = sqlite3.connect(DB) cleanup.execute(""" UPDATE recordings SET status = 'NEW' WHERE status = 'TRANSCRIBING' """) cleanup.commit() cleanup.close() raise elapsed = time.time() - started con = sqlite3.connect(DB) print() print("=" * 72) print("WHISPER-LAUF BEENDET") print("=" * 72) for status, count in con.execute(""" SELECT status, COUNT(*) FROM recordings GROUP BY status ORDER BY status """): print( f"{status:24} {count}" ) print( f"\nLaufzeit: " f"{elapsed / 60:.1f} Minuten" ) con.close() if __name__ == "__main__": mp.set_start_method( "spawn", force=True, ) main()