| """v3: пересборка 55-с групп с ЕСТЕСТВЕННЫМИ стыками (лечит «обрыв на точке», |
| константные паузы и мёртвые нули между клипами). |
| |
| Отличия от build_groups_ft4 (v2): |
| 1. Паузы на стыках — семплируются из ЭМПИРИЧЕСКИХ распределений реальных пауз |
| корпуса (v3/pause_hist.json, посчитан по 1.14M по-клиповых MFA-выравниваний), |
| класс — по хвостовому знаку клипа: медианы ~0.39s после '.', 0.42s после '?', |
| 0.27s после ',' (в v2 были константы 160/100 мс — в 2.5 раза короче естественных). |
| 2. Фейды адаптивные: до 40 мс приподнятым косинусом, но ТОЛЬКО внутри собственной |
| краевой тишины клипа (речь не трогаем; минимум 2 мс от щелчков). |
| 3. Паузы и хвост группы — не цифровой ноль, а КОМНАТНЫЙ ТОН клипа (самое тихое |
| 60-мс окно, тайлится зеркально — C0-непрерывно, с джиттером амплитуды). |
| 4. Клипы без MFA-выравнивания отсекаются на этапе плана (в v2 группы с ними |
| молча выпадали на этапе npy целиком). |
| 5. Оверсемпл вопросов (x2 прохода) и короткие группы (10%, 5-35 c) — встроены |
| (в v2 это делал отдельный extend_groups_e). |
| |
| Порядок клипов внутри спикера = порядок манифеста (для аудиокниг/подкастов это |
| порядок записи -> соседние клипы часто из одной сессии). |
| |
| Выход: ft4/groups_v3.json, wav в voxtream_ru/_groups55_v3/, grouped_v3_chunk*.parquet. |
| |
| Запуск: |
| .venv/bin/python ru/pipeline/build_groups_v3.py --plan-only # только план |
| .venv/bin/python ru/pipeline/build_groups_v3.py --workers 24 # план + wav |
| """ |
|
|
| import argparse |
| import json |
| import multiprocessing as mp |
| import random |
| import re |
| from pathlib import Path |
|
|
| import numpy as np |
| import pandas as pd |
| import soundfile as sf |
|
|
| ROOT = Path("/mnt/data/voxtream/ru_finetune/ru/data") |
| MFA_OUT = ROOT / "mfa_out" |
| GROUP_DIR = Path("/mnt/data/audio_data/voxtream_ru/_groups55_v3") |
| SR = 24000 |
| TARGET = int(54.5 * SR) |
| MIN_GROUP = int(35.0 * SR) |
| MIN_UNIQUE = int(15.0 * SR) |
| FADE_MIN = int(0.002 * SR) |
| FADE_MAX = int(0.040 * SR) |
| TONE_WIN = int(0.060 * SR) |
| LATIN = re.compile(r"[a-zA-Z]") |
|
|
| SHORT_MIN, SHORT_MAX = int(5 * SR), int(35 * SR) |
| Q_EXTRA_PASSES = 2 |
| SHORT_FRAC = 0.10 |
| Q_MIN_GROUP = int(10 * SR) |
|
|
| |
| |
| |
| GAP_CLAMP = {".": (0.10, 0.90), "!": (0.10, 1.05), "?": (0.10, 1.00), |
| ",": (0.06, 0.70), "none": (0.05, 0.60)} |
|
|
| OVERWRITE = False |
|
|
| |
| |
| |
| LUFS_ON = False |
| TARGET_LUFS = -23.0 |
| _meter = None |
|
|
|
|
| def lufs_gain(wav: np.ndarray) -> float: |
| global _meter |
| if len(wav) < SR // 2: |
| return 1.0 |
| if _meter is None: |
| import pyloudnorm |
| _meter = pyloudnorm.Meter(SR) |
| try: |
| loud = _meter.integrated_loudness(wav) |
| except Exception: |
| return 1.0 |
| if not np.isfinite(loud) or loud < -70: |
| return 1.0 |
| g = float(10 ** ((TARGET_LUFS - loud) / 20)) |
| peak = float(np.abs(wav).max()) * g |
| if peak > 0.99: |
| g *= 0.99 / peak |
| return min(max(g, 0.05), 20.0) |
|
|
|
|
| _hist = None |
|
|
|
|
| def load_hist(): |
| global _hist |
| h = json.load(open(ROOT / "v3" / "pause_hist.json")) |
| _hist = {k: np.asarray(v["sample"], dtype=np.float32) for k, v in h.items()} |
| return _hist |
|
|
|
|
| |
| GAP_MED = {".": 0.39, "!": 0.41, "?": 0.42, ",": 0.27, "none": 0.11} |
| |
| |
| SIGMA_WITHIN = 0.50 |
| |
| SIGMA_BETWEEN = 0.35 |
|
|
|
|
| def gap_of(text: str, rng: random.Random, tempo: float = 1.0) -> int: |
| """Пауза после клипа = медиана класса * ТЕМП ГРУППЫ * внутренний шум. |
| |
| v3 брал семпл из ОБЩЕКОРПУСНОЙ эмпирики (CV 0.70) на каждый стык |
| независимо — внутри одной группы возникал разброс, который в реальности |
| бывает только МЕЖДУ дикторами; модель выучила «после точки может быть что |
| угодно» и на инференсе давала «то слишком коротко, то слишком долго» |
| (жалоба пользователя; замер синтеза: медиана паузы 0.85с при норме 0.39с). |
| v4 расщепляет дисперсию: tempo — один множитель на группу (междикторская |
| компонента), шум CV 0.54 — внутридикторская. Суммарно даёт корпусную |
| вариативность, но СТРУКТУРИРОВАННУЮ. |
| """ |
| tail = str(text).rstrip()[-1:] |
| cls = tail if tail in ".!?," else "none" |
| lo, hi = GAP_CLAMP[cls] |
| g = GAP_MED[cls] * tempo * rng.lognormvariate(0.0, SIGMA_WITHIN) |
| return int(min(max(g, lo), hi) * SR) |
|
|
|
|
| def group_tempo(rng: random.Random) -> float: |
| """Логнормальный множитель темпа пауз группы (медиана 1.0, σ=0.35): |
| ~80% групп в диапазоне 0.64-1.57x — межспикерная компонента.""" |
| return float(min(max(rng.lognormvariate(0.0, SIGMA_BETWEEN), 0.5), 2.0)) |
|
|
|
|
| def aligned_indices() -> set: |
| """Индексы манифеста, у которых есть TextGrid (скан mfa_out, ~1 мин).""" |
| idx = set() |
| for d in MFA_OUT.iterdir(): |
| if not d.is_dir(): |
| continue |
| for f in d.iterdir(): |
| n = f.name |
| if n.endswith(".TextGrid"): |
| try: |
| idx.add(int(n[:-9])) |
| except ValueError: |
| pass |
| return idx |
|
|
|
|
| |
| def _rescue(short, pool, rng): |
| rescued, dropped = [], 0 |
| uniq_total = sum(p[3] for p in pool) |
| for g in short: |
| if uniq_total < MIN_UNIQUE: |
| dropped += 1 |
| continue |
| cand = [p for p in pool if p[0] not in set(g["idx"])] or list(pool) |
| rng.shuffle(cand) |
| ci = 0 |
| while g["total"] < MIN_GROUP and ci < 4 * len(cand): |
| idx, path, samples, slot = cand[ci % len(cand)] |
| ci += 1 |
| if g["total"] + slot > TARGET: |
| continue |
| for k, v in (("idx", idx), ("paths", path), ("samples", samples), ("slots", slot)): |
| g[k].append(v) |
| g["total"] += slot |
| g["reused"] = g.get("reused", 0) + 1 |
| if g["total"] >= MIN_GROUP: |
| rescued.append(g) |
| else: |
| dropped += 1 |
| return rescued, dropped |
|
|
|
|
| def _new(spk, kind="base"): |
| return {"speaker": spk, "idx": [], "paths": [], "slots": [], "samples": [], |
| "total": 0, "kind": kind} |
|
|
|
|
| def _push(g, idx, path, samples, slot): |
| g["idx"].append(idx) |
| g["paths"].append(path) |
| g["samples"].append(samples) |
| g["slots"].append(slot) |
| g["total"] += slot |
|
|
|
|
| def retempo(groups, texts, seed: int = 4242): |
| """v4: назначить каждой группе общий множитель темпа пауз и пересчитать слоты. |
| |
| Паузы внутри группы становятся согласованными (как у одного диктора в одной |
| сессии), а не независимыми выбросами из широкого распределения. Если после |
| пересчёта группа не влезает в TARGET — темп ужимается, в крайнем случае |
| дропается хвостовой клип. |
| """ |
| rng = random.Random(seed) |
| n_trim = 0 |
| for g in groups: |
| tempo = group_tempo(rng) |
| for _ in range(6): |
| slots = [smp + gap_of(texts[i], rng, tempo) if i >= 0 else slot |
| for i, smp, slot in zip(g["idx"], g["samples"], g["slots"])] |
| if sum(slots) <= TARGET: |
| break |
| tempo *= 0.85 |
| while sum(slots) > TARGET and len(slots) > 1: |
| for k in ("idx", "paths", "samples"): |
| g[k].pop() |
| slots.pop() |
| n_trim += 1 |
| g["slots"] = slots |
| g["total"] = sum(slots) |
| g["tempo"] = round(tempo, 3) |
| print(f"retempo: групп {len(groups)}, обрезано хвостовых клипов {n_trim}") |
| return [g for g in groups if g["total"] >= 5 * SR and g["idx"]] |
|
|
|
|
| def tempo_key(path: str) -> str: |
| """v6: версия темпа клипа (orig / slow / fast). |
| |
| Замер на v5: модель почти не копирует темп промпта (наклон 0.34). Одна из |
| причин — 18.8% групп СМЕШИВАЛИ оригиналы с темпо-аугментированными копиями: |
| промпт мог быть 0.75x, а продолжение 1.3x, т.е. данные прямо учили, что |
| темп промпта НЕ предсказывает темп речи. Группируем по (спикер, версия) — |
| внутри группы темп однороден, связь промпт->продолжение становится честной. |
| """ |
| if "_tempo_aug" not in path: |
| return "orig" |
| return "slow" if "_slow" in path else "fast" |
|
|
|
|
| |
| |
| |
| |
| TEMPO_EDGES = (4.3, 5.2, 6.1) |
|
|
|
|
| def tempo_bucket(sps: float) -> str: |
| if not np.isfinite(sps): |
| return "na" |
| return str(int(np.searchsorted(TEMPO_EDGES, sps))) |
|
|
|
|
| def pack_groups(df, aligned: set, seed: int = 42, clip_sps=None): |
| rng = random.Random(seed) |
| groups, n_lat, n_noal, n_rescued, n_dropped, n_reused = [], 0, 0, 0, 0, 0 |
| texts = df.text.astype(str) |
| by_spk_pool = {} |
|
|
| key = df.speaker.astype(str) + "|" + df.audio_path.map(tempo_key) |
| if clip_sps is not None: |
| key = key + "|" + pd.Series( |
| [tempo_bucket(s) for s in clip_sps[: len(df)]], index=df.index |
| ) |
| df = df.assign(_spk_tempo=key) |
| for spk_t, sub in df.groupby("_spk_tempo", sort=False): |
| spk = spk_t.split("|")[0] |
| cur, short, pool = None, [], [] |
| for row in sub.itertuples(): |
| if int(row.Index) not in aligned: |
| n_noal += 1 |
| continue |
| if LATIN.search(row.text): |
| n_lat += 1 |
| continue |
| samples = int(round(row.duration * SR)) |
| slot = samples + gap_of(row.text, rng) |
| if samples <= 0 or slot > TARGET: |
| continue |
| pool.append((int(row.Index), row.audio_path, samples, slot)) |
| if cur is None or cur["total"] + slot > TARGET: |
| if cur: |
| (groups if cur["total"] >= MIN_GROUP else short).append(cur) |
| cur = _new(spk) |
| _push(cur, int(row.Index), row.audio_path, samples, slot) |
| if cur: |
| (groups if cur["total"] >= MIN_GROUP else short).append(cur) |
| if short: |
| rescued, dropped = _rescue(short, pool, rng) |
| n_rescued += len(rescued) |
| n_dropped += dropped |
| n_reused += sum(g.get("reused", 0) for g in rescued) |
| groups.extend(rescued) |
| by_spk_pool[spk_t] = pool |
| n_base = len(groups) |
| print(f"базовых групп: {n_base}; латиница={n_lat}, без выравнивания={n_noal}, " |
| f"спасено={n_rescued} (+{n_reused} переисп.), дропнуто={n_dropped}") |
|
|
| |
| n_q = 0 |
| for spk_t, pool in by_spk_pool.items(): |
| spk = spk_t.rsplit("|", 1)[0] |
| qs = [p for p in pool if "?" in texts[p[0]]] |
| if not qs: |
| continue |
| for _ in range(Q_EXTRA_PASSES): |
| order = qs[:] |
| rng.shuffle(order) |
| cur = None |
| for idx, path, samples, slot in order: |
| if cur is None or cur["total"] + slot > TARGET: |
| if cur and cur["total"] >= Q_MIN_GROUP: |
| groups.append(cur) |
| cur = _new(spk, "question") |
| _push(cur, idx, path, samples, slot) |
| if cur and cur["total"] >= Q_MIN_GROUP: |
| groups.append(cur) |
| n_q = len(groups) - n_base |
|
|
| |
| spk_list = [s for s, p in by_spk_pool.items() if p] |
| n_short_target = int(SHORT_FRAC * n_base) |
| for _ in range(n_short_target): |
| spk_t = rng.choice(spk_list) |
| pool = by_spk_pool[spk_t] |
| cur = _new(spk_t.rsplit("|", 1)[0], "short") |
| target_len = rng.randint(SHORT_MIN, SHORT_MAX) |
| for _try in range(6): |
| idx, path, samples, slot = pool[rng.randrange(len(pool))] |
| if cur["total"] + slot > min(target_len + 3 * SR, TARGET): |
| if cur["total"] >= SHORT_MIN: |
| break |
| continue |
| _push(cur, idx, path, samples, slot) |
| if cur["total"] >= target_len: |
| break |
| if cur["idx"]: |
| groups.append(cur) |
| print(f"вопросных групп: +{n_q}, коротких: +{len(groups) - n_base - n_q}") |
| return retempo(groups, texts) |
|
|
|
|
| |
| def _cos_ramp(n: int) -> np.ndarray: |
| return (0.5 - 0.5 * np.cos(np.linspace(0, np.pi, n, dtype=np.float32))) |
|
|
|
|
| def adaptive_fades(wav: np.ndarray) -> np.ndarray: |
| """Фейды внутри краевой тишины клипа: до 40 мс, речь не трогаем.""" |
| peak = float(np.max(np.abs(wav))) + 1e-9 |
| loud = np.abs(wav) > max(0.02 * peak, 1e-4) |
| if not loud.any(): |
| return wav |
| first = int(np.argmax(loud)) |
| last = len(wav) - int(np.argmax(loud[::-1])) |
| fi = min(max(first, FADE_MIN), FADE_MAX, len(wav)) |
| fo = min(max(len(wav) - last, FADE_MIN), FADE_MAX, len(wav)) |
| wav[:fi] *= _cos_ramp(fi) |
| wav[-fo:] *= _cos_ramp(fo)[::-1] |
| return wav |
|
|
|
|
| def room_tone(wav: np.ndarray) -> np.ndarray: |
| """Самое тихое 120-мс окно клипа (шаблон спектра); тихих нет — глушим до -46 dBFS.""" |
| win = 2 * TONE_WIN |
| if len(wav) < 2 * win: |
| return np.zeros(win, dtype=np.float32) |
| hop = win // 2 |
| n = (len(wav) - win) // hop |
| view = np.lib.stride_tricks.sliding_window_view(wav, win)[::hop][:n] |
| rms = np.sqrt((view ** 2).mean(axis=1)) |
| k = int(np.argmin(rms)) |
| tone = view[k].copy() |
| if rms[k] > 5e-3: |
| tone *= 5e-3 / rms[k] |
| return tone |
|
|
|
|
| _N_FFT = 512 |
| _HOP = 256 |
| _HANN = np.hanning(_N_FFT).astype(np.float32) |
|
|
|
|
| def synth_tone(template: np.ndarray, length: int, seed: int) -> np.ndarray: |
| """Стационарный шум со спектральной огибающей шаблона — БЕЗ зацикливания. |
| |
| v3 первой версии тайлил 60-мс окно зеркально: период ~8 Гц слышен как |
| «вертолёт», модель выучила текстуру (жалоба юзера). Здесь — средняя |
| STFT-магнитуда шаблона + случайные фазы на каждый кадр + overlap-add: |
| цвет комнаты сохранён, повторов нет в принципе. |
| """ |
| if length <= 0: |
| return np.zeros(0, dtype=np.float32) |
| rng = np.random.default_rng(seed) |
| if len(template) < _N_FFT or not np.any(template): |
| return np.zeros(length, dtype=np.float32) |
| n = (len(template) - _N_FFT) // _HOP + 1 |
| frames = np.lib.stride_tricks.sliding_window_view(template, _N_FFT)[::_HOP][:n] |
| mag = np.abs(np.fft.rfft(frames * _HANN, axis=1)).mean(axis=0) |
|
|
| n_frames = length // _HOP + 3 |
| phases = rng.uniform(0, 2 * np.pi, size=(n_frames, len(mag))) |
| sig_frames = np.fft.irfft(mag * np.exp(1j * phases), n=_N_FFT, axis=1).real |
| sig_frames *= _HANN |
| out = np.zeros(n_frames * _HOP + _N_FFT, dtype=np.float32) |
| for i in range(n_frames): |
| out[i * _HOP:i * _HOP + _N_FFT] += sig_frames[i] |
| out = out[_N_FFT: _N_FFT + length] |
| t_rms = float(np.sqrt((template ** 2).mean())) |
| o_rms = float(np.sqrt((out ** 2).mean())) + 1e-12 |
| return (out * (t_rms / o_rms)).astype(np.float32) |
|
|
|
|
| def fill_tone(buf: np.ndarray, start: int, end: int, tone: np.ndarray, seed: int): |
| """Заполняет [start:end) синтезированным комнатным тоном с 5-мс рампами.""" |
| if end <= start or not len(tone): |
| return |
| buf[start:end] = synth_tone(tone, end - start, seed) |
| r = min(int(0.005 * SR), (end - start) // 2) |
| if r > 0: |
| buf[start:start + r] *= _cos_ramp(r) |
| buf[end - r:end] *= _cos_ramp(r)[::-1] |
|
|
|
|
| def write_group(task): |
| gi, g = task |
| out = GROUP_DIR / f"g{gi:06d}.wav" |
| if not OVERWRITE and out.exists(): |
| try: |
| if sf.info(str(out)).frames == int(55 * SR): |
| return gi, str(out), "skip" |
| except Exception: |
| out.unlink(missing_ok=True) |
| buf = np.zeros(int(55 * SR), dtype=np.float32) |
| pos = 0 |
| try: |
| tone = np.zeros(2 * TONE_WIN, dtype=np.float32) |
| for j, (path, samples, slot) in enumerate(zip(g["paths"], g["samples"], g["slots"])): |
| wav, sr = sf.read(path, dtype="float32") |
| assert sr == SR, f"{path}: sr={sr}" |
| if wav.ndim > 1: |
| wav = wav.mean(axis=1) |
| wav = wav[:samples].copy() |
| if LUFS_ON: |
| wav *= lufs_gain(wav) |
| wav = adaptive_fades(wav) |
| buf[pos:pos + len(wav)] = wav |
| tone = room_tone(wav) |
| fill_tone(buf, pos + len(wav), pos + slot, tone, 100_000 + gi * 64 + j) |
| pos += slot |
| fill_tone(buf, pos, len(buf), tone, 100_000 + gi * 64 + 63) |
| sf.write(out, buf, SR, subtype="PCM_16") |
| return gi, str(out), "ok" |
| except Exception as e: |
| out.unlink(missing_ok=True) |
| return gi, "", f"fail:{e}" |
|
|
|
|
| def main(): |
| ap = argparse.ArgumentParser() |
| ap.add_argument("--plan-only", action="store_true") |
| ap.add_argument("--workers", type=int, default=24) |
| ap.add_argument("--chunk-size", type=int, default=70_000) |
| ap.add_argument("--overwrite", action="store_true", |
| help="перерендерить wav даже если файл на месте (план не меняется)") |
| global OVERWRITE, MFA_OUT, GROUP_DIR |
| ap.add_argument("--mfa-out", default="mfa_out", help="каталог TextGrid (v4: mfa_out_v4)") |
| ap.add_argument("--out-groups", default="groups_v3.json") |
| ap.add_argument("--group-dir", default="/mnt/data/audio_data/voxtream_ru/_groups55_v3") |
| ap.add_argument("--parquet-prefix", default="grouped_v3_chunk") |
| ap.add_argument("--tempo-buckets", action="store_true", |
| help="v7: группы однородны по темпу (clip_sps.npy)") |
| ap.add_argument("--clip-sps", default="clip_sps.npy", |
| help="файл SPS по клипам для темпо-корзин") |
| ap.add_argument("--lufs", action="store_true", |
| help="v10: каждый клип -> -23 LUFS при рендере") |
| args = ap.parse_args() |
| global LUFS_ON |
| OVERWRITE = args.overwrite |
| LUFS_ON = args.lufs |
| MFA_OUT = ROOT / args.mfa_out |
| GROUP_DIR = Path(args.group_dir) |
| if LUFS_ON: |
| print(f"LUFS-нормализация ВКЛ: клипы -> {TARGET_LUFS} LUFS") |
|
|
| load_hist() |
| print("скан выравниваний…", flush=True) |
| aligned = aligned_indices() |
| print(f"выровненных клипов: {len(aligned)}", flush=True) |
|
|
| df = pd.read_csv(ROOT / "manifest.csv", sep="|", low_memory=False) |
| clip_sps = None |
| if args.tempo_buckets: |
| clip_sps = np.load(ROOT / args.clip_sps) |
| print(f"темпо-корзины ВКЛ: границы {TEMPO_EDGES} слог/с") |
| groups = pack_groups(df, aligned, clip_sps=clip_sps) |
| total_h = sum(g["total"] for g in groups) / SR / 3600 |
| used = len({i for g in groups for i in g["idx"]}) |
| print(f"групп: {len(groups)}, уникальных клипов: {used}/{len(df)}, ~{total_h:.1f} ч слотов") |
| json.dump(groups, open(ROOT / args.out_groups, "w"), ensure_ascii=False) |
| if args.plan_only: |
| return |
|
|
| GROUP_DIR.mkdir(parents=True, exist_ok=True) |
| fails = 0 |
| with mp.Pool(args.workers) as pool: |
| for gi, path, status in pool.imap(write_group, enumerate(groups), chunksize=16): |
| groups[gi]["group_wav"] = path |
| fails += status.startswith("fail") |
| if (gi + 1) % 10_000 == 0: |
| print(f" {gi + 1}/{len(groups)} (fail={fails})", flush=True) |
| groups = [g for g in groups if g.get("group_wav")] |
| json.dump(groups, open(ROOT / args.out_groups, "w"), ensure_ascii=False) |
| print(f"записано групп: {len(groups)} (fail={fails})") |
|
|
| paths = [g["group_wav"] for g in groups] |
| for ci in range(0, len(paths), args.chunk_size): |
| p = ROOT / f"{args.parquet_prefix}{ci // args.chunk_size}.parquet" |
| pd.DataFrame({"paths": paths[ci:ci + args.chunk_size]}).to_parquet(p, index=False) |
| print(f"{p.name}: {min(args.chunk_size, len(paths) - ci)} путей") |
|
|
|
|
| if __name__ == "__main__": |
| main() |
|
|