LUNA / Base /scripts /consolidate_datasets.py
ASTERIZER
LUNA 100M: cloud-ready training pipeline
ad68b7f
Raw
History Blame Contribute Delete
8.55 kB
"""
Consolidate all pretraining and finetuning datasets into single files.
Creates:
Base/data/consolidated_pretrain/ -> all_pretrain_data.parquet (3B + English)
Base/Datasets/consolidated_finetune/ -> all_finetune_data.json (v1 + English instruct)
Also counts tokens using the Pythia tokenizer (batch mode for speed).
"""
import os
import json
import glob
import pandas as pd
from pathlib import Path
from tokenizers import Tokenizer
ROOT = Path(__file__).resolve().parent.parent.parent # LUNA root
# ── Tokenizer ──────────────────────────────────────────────────────
print("Loading tokenizer...")
tokenizer = Tokenizer.from_file(
str(ROOT / "Base" / "checkpoints" / "EleutherAI" / "pythia-160m" / "tokenizer.json")
)
BATCH_SIZE = 10000 # encode_batch processes this many at once
def count_tokens_text(texts):
"""Count total tokens using fast batch encoding."""
total = 0
for i in range(0, len(texts), BATCH_SIZE):
batch = texts[i : i + BATCH_SIZE]
encoded = tokenizer.encode_batch(batch, add_special_tokens=False)
total += sum(len(e.ids) for e in encoded)
done = min(i + BATCH_SIZE, len(texts))
if done % 100000 == 0 or done == len(texts):
print(f" Tokenized {done:,}/{len(texts):,} documents ({total:,} tokens)")
return total
# ══════════════════════════════════════════════════════════════════
# 1. PRETRAINING DATA (filtered_3b + filtered_english)
# ══════════════════════════════════════════════════════════════════
print("\n" + "=" * 60)
print("CONSOLIDATING PRETRAINING DATA")
print("=" * 60)
pretrain_out = ROOT / "Base" / "data" / "consolidated_pretrain"
pretrain_out.mkdir(parents=True, exist_ok=True)
# Read filtered_3b (15 parquets)
dir_3b = ROOT / "Base" / "data" / "filtered_3b"
files_3b = sorted(glob.glob(str(dir_3b / "*.parquet")))
print(f"\nReading filtered_3b: {len(files_3b)} files...")
dfs_3b = [pd.read_parquet(f) for f in files_3b]
df_3b = pd.concat(dfs_3b, ignore_index=True)
print(f" -> {len(df_3b):,} documents")
# Read filtered_english (fineweb + wiki parquets)
dir_en = ROOT / "Base" / "data" / "filtered_english"
files_en = sorted(glob.glob(str(dir_en / "*.parquet")))
print(f"\nReading filtered_english: {len(files_en)} files...")
dfs_en = [pd.read_parquet(f) for f in files_en]
df_en = pd.concat(dfs_en, ignore_index=True)
print(f" -> {len(df_en):,} documents")
# Add source labels
df_3b["source"] = "filtered_3b"
df_en_list = []
for f in files_en:
fname = os.path.basename(f)
tmp = pd.read_parquet(f)
if fname.startswith("fineweb"):
tmp["source"] = "english_fineweb"
elif fname.startswith("wiki"):
tmp["source"] = "english_wiki"
else:
tmp["source"] = "english_other"
df_en_list.append(tmp)
df_en_labeled = pd.concat(df_en_list, ignore_index=True)
# Combine all pretraining data
df_pretrain = pd.concat([df_3b, df_en_labeled], ignore_index=True)
out_path = pretrain_out / "all_pretrain_data.parquet"
df_pretrain.to_parquet(str(out_path), index=False)
print(f"\nSaved combined pretrain: {out_path}")
print(f" Total documents: {len(df_pretrain):,}")
# Token counting
print("\nCounting pretraining tokens (this may take a while)...")
pretrain_tokens_3b = count_tokens_text(df_3b["text"].tolist())
print(f" filtered_3b tokens: {pretrain_tokens_3b:,}")
pretrain_tokens_en = count_tokens_text(df_en["text"].tolist())
print(f" filtered_english tokens: {pretrain_tokens_en:,}")
pretrain_total = pretrain_tokens_3b + pretrain_tokens_en
print(f" TOTAL pretrain tokens: {pretrain_total:,}")
# ══════════════════════════════════════════════════════════════════
# 2. FINETUNING DATA (finetune + finetune_english, train+val each)
# ══════════════════════════════════════════════════════════════════
print("\n" + "=" * 60)
print("CONSOLIDATING FINETUNING DATA")
print("=" * 60)
finetune_out = ROOT / "Base" / "Datasets" / "consolidated_finetune"
finetune_out.mkdir(parents=True, exist_ok=True)
# Read v1 finetune
ft_v1_train = json.loads((ROOT / "Base" / "Datasets" / "finetune" / "train.json").read_text(encoding="utf-8"))
ft_v1_val = json.loads((ROOT / "Base" / "Datasets" / "finetune" / "val.json").read_text(encoding="utf-8"))
print(f"\nFinetune v1: train={len(ft_v1_train):,} val={len(ft_v1_val):,}")
# Read English finetune
ft_en_train = json.loads((ROOT / "Base" / "Datasets" / "finetune_english" / "train.json").read_text(encoding="utf-8"))
ft_en_val = json.loads((ROOT / "Base" / "Datasets" / "finetune_english" / "val.json").read_text(encoding="utf-8"))
print(f"Finetune EN: train={len(ft_en_train):,} val={len(ft_en_val):,}")
# Add source tags
for item in ft_v1_train + ft_v1_val:
item["source"] = "finetune_v1"
for item in ft_en_train + ft_en_val:
item["source"] = "finetune_english"
# Combine
all_finetune = ft_v1_train + ft_v1_val + ft_en_train + ft_en_val
out_path_ft = finetune_out / "all_finetune_data.json"
with open(out_path_ft, "w", encoding="utf-8") as f:
json.dump(all_finetune, f, ensure_ascii=False, indent=2)
print(f"\nSaved combined finetune: {out_path_ft}")
print(f" Total samples: {len(all_finetune):,}")
# Token counting for finetuning
print("\nCounting finetuning tokens...")
def finetune_text(item):
"""Reconstruct the full text that gets tokenized during finetuning."""
parts = []
if item.get("instruction"):
parts.append(item["instruction"])
if item.get("input"):
parts.append(item["input"])
if item.get("output"):
parts.append(item["output"])
return " ".join(parts)
ft_v1_texts = [finetune_text(x) for x in ft_v1_train + ft_v1_val]
ft_en_texts = [finetune_text(x) for x in ft_en_train + ft_en_val]
ft_v1_tokens = count_tokens_text(ft_v1_texts)
print(f" finetune_v1 tokens: {ft_v1_tokens:,}")
ft_en_tokens = count_tokens_text(ft_en_texts)
print(f" finetune_english tokens: {ft_en_tokens:,}")
ft_total = ft_v1_tokens + ft_en_tokens
print(f" TOTAL finetune tokens: {ft_total:,}")
# ══════════════════════════════════════════════════════════════════
# SUMMARY
# ══════════════════════════════════════════════════════════════════
print("\n" + "=" * 60)
print("FINAL SUMMARY")
print("=" * 60)
summary = f"""
PRETRAINING DATA (Base/data/consolidated_pretrain/all_pretrain_data.parquet)
Source: filtered_3b -> {len(df_3b):>10,} docs | {pretrain_tokens_3b:>15,} tokens
Source: filtered_english -> {len(df_en):>10,} docs | {pretrain_tokens_en:>15,} tokens
─────────────────────────────────────────────────────────
TOTAL -> {len(df_pretrain):>10,} docs | {pretrain_total:>15,} tokens
FINETUNING DATA (Base/Datasets/consolidated_finetune/all_finetune_data.json)
Source: finetune_v1 -> {len(ft_v1_train)+len(ft_v1_val):>10,} samples | {ft_v1_tokens:>15,} tokens
Source: finetune_english -> {len(ft_en_train)+len(ft_en_val):>10,} samples | {ft_en_tokens:>15,} tokens
─────────────────────────────────────────────────────────
TOTAL -> {len(all_finetune):>10,} samples | {ft_total:>15,} tokens
GRAND TOTAL TOKENS: {pretrain_total + ft_total:,}
"""
print(summary)
# Save summary as text file too
with open(pretrain_out.parent.parent / "data" / "consolidated_pretrain" / "SUMMARY.txt", "w") as f:
f.write(summary)
with open(finetune_out / "SUMMARY.txt", "w") as f:
f.write(summary)
print("Summary saved to both consolidated folders.")