"""Processing data for pretraining.""" import argparse import os import sys _PROJECT_ROOT = os.path.abspath(os.path.dirname(__file__)) while _PROJECT_ROOT and not os.path.isdir(os.path.join(_PROJECT_ROOT, "model")): _PARENT = os.path.dirname(_PROJECT_ROOT) if _PARENT == _PROJECT_ROOT: break _PROJECT_ROOT = _PARENT _MODEL_ROOT = os.path.join(_PROJECT_ROOT, "model") _ONESCIENCE_ROOT = os.environ.get("ONESCIENCE_ROOT") for _path in (_MODEL_ROOT, _PROJECT_ROOT): if os.path.exists(_path) and _path not in sys.path: sys.path.insert(0, _path) if _ONESCIENCE_ROOT: _ONESCIENCE_SRC = os.path.join(_ONESCIENCE_ROOT, "src") for _path in (_ONESCIENCE_SRC, _ONESCIENCE_ROOT): if os.path.exists(_path) and _path not in sys.path: sys.path.insert(0, _path) import multiprocessing import os import sys import lm_dataformat as lmd import numpy as np sys.path.append( os.path.abspath(os.path.join(os.path.dirname(__file__), os.path.pardir)) ) import time from abc import ABC, abstractmethod from threading import Semaphore from typing import List, Union import ftfy import tqdm # from evo2.tokenizer import build_tokenizer from evo2.data import indexed_dataset def build_tokenizer(args): """Initialize tokenizer.""" if args.rank == 0: print("> building {} tokenizer ...".format(args.tokenizer_type), flush=True) # Select and instantiate the tokenizer. if args.tokenizer_type.lower() == "CharLevelTokenizer".lower(): tokenizer = CharLevelTokenizer(vocab_size=512) else: raise NotImplementedError( "{} tokenizer is not " "implemented.".format(args.tokenizer_type) ) # Add vocab size. args.padded_vocab_size = _vocab_size_with_padding(tokenizer.vocab_size, args) return tokenizer def _vocab_size_with_padding(orig_vocab_size, args): """Pad vocab size so it is divisible by model parallel size and still having GPU friendly size.""" after = orig_vocab_size multiple = args.make_vocab_size_divisible_by * args.model_parallel_size while (after % multiple) != 0: after += 1 if args.rank == 0: print( " > padded vocab (size: {}) with {} dummy tokens " "(new size: {})".format(orig_vocab_size, after - orig_vocab_size, after), flush=True, ) return after class AbstractTokenizer(ABC): """Abstract class for tokenizer.""" def __init__(self, name): self.name = name super().__init__() @property @abstractmethod def vocab_size(self): pass @property @abstractmethod def vocab(self): """Dictionary from vocab text token to id token.""" @property @abstractmethod def inv_vocab(self): """Dictionary from vocab id token to text token.""" @abstractmethod def tokenize(self, text): pass def detokenize(self, token_ids): raise NotImplementedError( "detokenizer is not implemented for {} " "tokenizer".format(self.name) ) @property def cls(self): raise NotImplementedError( "CLS is not provided for {} " "tokenizer".format(self.name) ) @property def sep(self): raise NotImplementedError( "SEP is not provided for {} " "tokenizer".format(self.name) ) @property def pad(self): raise NotImplementedError( "PAD is not provided for {} " "tokenizer".format(self.name) ) @property def eod(self): raise NotImplementedError( "EOD is not provided for {} " "tokenizer".format(self.name) ) @property def mask(self): raise NotImplementedError( "MASK is not provided for {} " "tokenizer".format(self.name) ) class CharLevelTokenizer(AbstractTokenizer): """Character Level Tokenizer""" def __init__(self, vocab_size): name = "CharLevelTokenizer" super().__init__(name) self._vocab_size = vocab_size self.eod_id = 0 self.pad_id = 1 self._used_tokens = set() # 璁板綍鐢ㄨ繃鐨?token id def clamp(self, n): return max(32, min(n, self.vocab_size)) @property def vocab_size(self): return self._vocab_size @property def vocab(self): raise NotImplementedError @property def inv_vocab(self): raise NotImplementedError def decode_token(self, token: int): return str(chr(self.clamp(token))) # def tokenize(self, text: str): # return list(np.fromstring(text, dtype=np.uint8)) def tokenize(self, text: str): tokens = list(np.fromstring(text, dtype=np.uint8)) # 璁板綍鎵€鏈?clamp 鍚庣殑 token id clamped_tokens = [self.clamp(t) for t in tokens] self._used_tokens.update(clamped_tokens) return tokens def tokenize_batch(self, text_batch: Union[List[str], str]): if isinstance(text_batch, list): return [self.tokenize(s) for s in text_batch] else: return self.tokenize(text_batch) def detokenize(self, token_ids): return "".join(list(map(self.decode_token, token_ids))) @property def eod(self): return self.eod_id @property def pad(self): return self.pad_id class Encoder(object): def __init__(self, args): self.args = args def initializer(self): # Use Encoder class as a container for global data Encoder.tokenizer = build_tokenizer(self.args) def encode(self, text): if self.args.ftfy: text = ftfy.fix_text(text) ids = {} for key in self.args.jsonl_keys: doc_ids = [] text_ids = Encoder.tokenizer.tokenize(text) if ( self.args.enforce_sample_length and (len(text_ids) + int(self.args.append_eod)) > self.args.enforce_sample_length ): raise ValueError( "Detected input text with a length greater than the maximum " f"possible sample length of {self.args.enforce_sample_length}.)" ) if len(text_ids) > 0: doc_ids.append(text_ids) if self.args.append_eod: doc_ids[-1].append(Encoder.tokenizer.eod) if self.args.enforce_sample_length: # Pad up to max sequence length. doc_ids[-1] += [Encoder.tokenizer.pad] * ( self.args.enforce_sample_length - len(doc_ids[-1]) ) ids[key] = doc_ids return ids, len(text) def get_args(): parser = argparse.ArgumentParser() group = parser.add_argument_group(title="input data") group.add_argument( "--input", type=str, required=True, help="Path to input jsonl files or lmd archive(s) - if using multiple archives, put them in a comma separated " "list", ) group.add_argument( "--jsonl-keys", nargs="+", default=["text"], help="space separate listed of keys to extract from jsonl. Defa", ) group.add_argument( "--num-docs", default=None, help="Optional: Number of documents in the input data (if known) for an accurate progress bar.", type=int, ) group = parser.add_argument_group(title="tokenizer") group.add_argument( "--tokenizer-type", type=str, required=True, choices=[ "HFGPT2Tokenizer", "HFTokenizer", "GPT2BPETokenizer", "CharLevelTokenizer", "TiktokenTokenizer", ], help="What type of tokenizer to use.", ) group.add_argument( "--vocab-file", type=str, default=None, help="Path to the vocab file" ) group.add_argument( "--merge-file", type=str, default=None, help="Path to the BPE merge file (if necessary).", ) group.add_argument( "--append-eod", action="store_true", help="Append an token to the end of a document.", ) group.add_argument( "--enforce-sample-length", type=int, default=None, help="Forces all samples to have the specified length. If shorter, pads up to the length. If longer, throws an error.", ) group.add_argument("--ftfy", action="store_true", help="Use ftfy to clean text") group = parser.add_argument_group(title="output data") group.add_argument( "--output-prefix", type=str, required=True, help="Path to binary output file without suffix", ) group.add_argument( "--dataset-impl", type=str, default="mmap", choices=["lazy", "cached", "mmap"], help="Dataset implementation to use. Default: mmap", ) group = parser.add_argument_group(title="runtime") group.add_argument( "--workers", type=int, default=1, help="Number of worker processes to launch" ) group.add_argument( "--log-interval", type=int, default=100, help="Interval between progress updates", ) args = parser.parse_args() args.keep_empty = False # some default/dummy values for the tokenizer args.rank = 0 args.make_vocab_size_divisible_by = 128 args.model_parallel_size = 1 return args def yield_from_files(fnames: list, semaphore): """ Iterator over input documents using lm_dataformat. Should be able to handle jsons / texts / other compressed formats. Also filters out empty documents. :param fnames: list of filenames """ def yielder(fname, semaphore): for f in filter(lambda x: x, lmd.Reader(fname).stream_data()): semaphore.acquire() # import pdb;pdb.set_trace() yield f for fname in fnames: semaphore.acquire() yield from yielder(fname, semaphore) def main(): args = get_args() encoder = Encoder(args) tokenizer = build_tokenizer(args) print(f"Vocab size: {tokenizer.vocab_size}") print(f"Output prefix: {args.output_prefix}") # build a semaphore object to stop `yield_from_files` from getting ahead of encoder.encode and # hence building up memory semaphore = Semaphore(10000 + args.workers) # use multiprocessing to iterate over input documents # import pdb;pdb.set_trace() fin = yield_from_files(args.input.split(","), semaphore) print(fin) if args.workers > 1: pool = multiprocessing.Pool(args.workers, initializer=encoder.initializer) encoded_docs = pool.imap(encoder.encode, fin, chunksize=25) else: encoder.initializer() encoded_docs = (encoder.encode(doc) for doc in fin) # print(Encoder.tokenizer._used_tokens) # make a dataset builder for each key in args.jsonl_keys # each key will output to a different file beginning with args.output_prefix output_bin_files = {} output_idx_files = {} builders = {} tokenizer_name = tokenizer.name.replace(" ", "") for key in args.jsonl_keys: output_bin_files[key] = "{}_{}_{}_{}.bin".format( args.output_prefix, key, tokenizer_name, "document" ) output_idx_files[key] = "{}_{}_{}_{}.idx".format( args.output_prefix, key, tokenizer_name, "document" ) builders[key] = indexed_dataset.make_builder( output_bin_files[key], impl=args.dataset_impl, vocab_size=tokenizer.vocab_size, ) # actually do tokenization proc_start = time.time() total_bytes_processed = 0 pbar = tqdm.tqdm() for i, (doc, bytes_processed) in enumerate(encoded_docs, start=1): total_bytes_processed += bytes_processed # release semaphore so `yield_from_files` can add another file to the buffer semaphore.release() # add each tokenized document / sentence for key, sentences in doc.items(): for sentence in sentences: builders[key].add_item(np.array(sentence, dtype=builders[key].dtype)) # tell the builder that a document has finished builders[key].end_document() # log progress if i % args.log_interval == 0: current = time.time() elapsed = current - proc_start mbs = total_bytes_processed / elapsed / 1024 / 1024 pbar.set_description( f"Processed {i}{'' if args.num_docs is None else '/' + str(args.num_docs)} documents ({i / elapsed} docs/s, {mbs} MB/s)." ) if i != 0: pbar.update(args.log_interval) # save output file for key in args.jsonl_keys: builders[key].finalize(output_idx_files[key]) if __name__ == "__main__": main()