File size: 1,757 Bytes
d063ac1
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
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)