fractus-cte-atom / scripts /convert_atom.py
thefinalboss's picture
convert_atom.py skips existing.
d063ac1 verified
Raw History Blame Contribute Delete
1.76 kB
import sys, time, json
from pathlib import Path
from concurrent.futures import ProcessPoolExecutor, as_completed
import numpy as np
sys.path.insert(0, "/workspace/fractus-cte-atom-main")
from fractus.atom_tokenizer import AtomFractusTokenizer
OUT = Path("/workspace/atom_corpus")
OUT.mkdir(exist_ok=True)
def convert_one(path):
tok = AtomFractusTokenizer()
out_path = OUT / (Path(path).stem + ".i16")
if out_path.exists() and out_path.stat().st_size > 0:
return Path(path).name, 0, 0, True
ids = []
n = 0
with open(path) as f:
for line in f:
line = line.strip()
if not line:
continue
obj = json.loads(line)
text = "\n".join(m.get("content","") for m in obj.get("messages",[]) if isinstance(m, dict))
if not text.strip():
continue
enc, _ = tok.encode_with_features(text)
ids.extend(enc)
n += 1
if ids:
arr = np.asarray(ids, dtype=np.int16)
arr.tofile(out_path)
return Path(path).name, n, len(ids), False
files = sorted(Path("/workspace/fractus-datasets").rglob("*.jsonl"))
print(f"files {len(files)}", flush=True)
total = 0
skipped = 0
t0 = time.time()
with ProcessPoolExecutor(max_workers=8) as ex:
futs = {ex.submit(convert_one, str(p)): p for p in files}
for i, fut in enumerate(as_completed(futs)):
name, n, s, was_skip = fut.result()
if was_skip:
skipped += 1
else:
total += s
if i % 5 == 0:
print(f"{i}/{len(files)} {name} docs={n} ids={s} new={total:,} skipped={skipped} {total/max(time.time()-t0,1):.0f}/s", flush=True)
print(f"DONE new={total:,} skipped={skipped}", flush=True)