# Copyright (c) 2023 PaddlePaddle Authors. All Rights Reserved. # # Licensed under the Apache License, Version 2.0 (the "License"); # you may not use this file except in compliance with the License. # You may obtain a copy of the License at # # http://www.apache.org/licenses/LICENSE-2.0 # # Unless required by applicable law or agreed to in writing, software # distributed under the License is distributed on an "AS IS" BASIS, # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # See the License for the specific language governing permissions and # limitations under the License. # https://github.com/NVIDIA/Megatron-LM/blob/060415572f4365a2e895f8036c4e37dad0efbdf5/megatron/data/indexed_dataset.py # Copyright (c) Facebook, Inc. and its affiliates. # # This source code is licensed under the MIT license found in the # LICENSE file in the root directory of this source tree. # copied from fairseq/fairseq/data/indexed_dataset.py # Removed IndexedRawTextDataset since it relied on Fairseq dictionary # other slight modifications to remove fairseq dependencies # Added document index to index file and made it accessible. # An empty sentence no longer separates documents. import os import shutil import struct import time from dataclasses import fields from functools import lru_cache from itertools import accumulate import numpy as np import paddle def print_rank_0(*args, **kwargs): if paddle.distributed.get_rank() == 0: print(*args, **kwargs) def __best_fitting_dtype(vocab_size=None): if vocab_size is not None and vocab_size < 65500: return np.uint16 else: return np.int32 def get_available_dataset_impl(): return ["lazy", "mmap"] def make_dataset(path, impl, skip_warmup=False): if CompatibleIndexedDataset.exists(path): print("Using old dataset (.npy & .npz)") return CompatibleIndexedDataset(path) elif not IndexedDataset.exists(path): print(f"Dataset does not exist: {path}") print("Path should be a basename that both .idx and .bin can be appended to get full filenames.") return None elif impl == "lazy" and IndexedDataset.exists(path): return IndexedDataset(path) elif impl == "mmap" and MMapIndexedDataset.exists(path): return MMapIndexedDataset(path, skip_warmup) print(f"Unknown dataset implementation: {impl}") return None def make_sft_dataset(path, dataclass, skip_warmup=False, impl="mmap"): if impl != "mmap": raise ValueError("SFT Indexed Dataset only support mmap memory-mapped method temporarily") print_rank_0(" > building dataset index ...") start_time = time.time() sft_indexed_dataset = SFTMMapIndexedDataset(path, dataclass, skip_warmup) print_rank_0(" > finished creating SFT indexed dataset in {:4f} " "seconds".format(time.time() - start_time)) print_rank_0(" number of samples: {}".format(len(sft_indexed_dataset.doc_idx) - 1)) return sft_indexed_dataset def dataset_exists(path, impl): if impl == "mmap": return MMapIndexedDataset.exists(path) else: return IndexedDataset.exists(path) def read_longs(f, n): a = np.empty(n, dtype=np.int64) f.readinto(a) return a def write_longs(f, a): f.write(np.array(a, dtype=np.int64)) def read_shorts(f, n): a = np.empty(n, dtype=np.int32) f.readinto(a) return a def write_shorts(f, a): f.write(np.array(a, dtype=np.int32)) dtypes = { 1: np.uint8, 2: np.int8, 3: np.int16, 4: np.int32, 5: np.int64, 6: np.float64, 7: np.float32, 8: np.uint16, 9: np.uint32, 10: np.uint64, } def code(dtype): for k in dtypes.keys(): if dtypes[k] == dtype: return k raise ValueError(dtype) def index_file_path(prefix_path): return prefix_path + ".idx" def sft_index_file_path(prefix_path): return os.path.join(prefix_path, "index.idx") def sft_data_file_path(prefix_path, dataclass): file_path_list = [] for field in fields(dataclass): file_path = os.path.join(prefix_path, f"{field.name}.bin") file_path_list.append(file_path) return file_path_list def data_file_path(prefix_path): return prefix_path + ".bin" def loss_mask_file_path(prefix_path): return prefix_path + ".lsm" def create_doc_idx(sizes): doc_idx = [0] for i, s in enumerate(sizes): if s == 0: doc_idx.append(i + 1) return doc_idx class IndexedDataset(paddle.io.Dataset): """Loader for IndexedDataset""" _HDR_MAGIC = b"TNTIDX\x00\x00" def __init__(self, path): super().__init__() self.path = path self.data_file = None self.read_index(path) def read_index(self, path): with open(index_file_path(path), "rb") as f: magic = f.read(8) assert magic == self._HDR_MAGIC, ( "Index file doesn't match expected format. " "Make sure that --dataset-impl is configured properly." ) version = f.read(8) assert struct.unpack("= self._len: raise IndexError("index out of range") def __del__(self): if self.data_file: self.data_file.close() # @lru_cache(maxsize=8) def __getitem__(self, idx): if not self.data_file: self.read_data(self.path) if isinstance(idx, int): i = idx self.check_index(i) tensor_size = self.sizes[self.dim_offsets[i] : self.dim_offsets[i + 1]] a = np.empty(tensor_size, dtype=self.dtype) self.data_file.seek(self.data_offsets[i] * self.element_size) self.data_file.readinto(a) return a elif isinstance(idx, slice): start, stop, step = idx.indices(len(self)) if step != 1: raise ValueError("Slices into indexed_dataset must be contiguous") sizes = self.sizes[self.dim_offsets[start] : self.dim_offsets[stop]] size = sum(sizes) a = np.empty(size, dtype=self.dtype) self.data_file.seek(self.data_offsets[start] * self.element_size) self.data_file.readinto(a) offsets = list(accumulate(sizes)) sents = np.split(a, offsets[:-1]) return sents def get(self, idx, offset=0, length=None): """Retrieves a single item from the dataset with the option to only return a portion of the item. get(idx) is the same as [idx] but get() does not support slicing. """ if not self.data_file: self.read_data(self.path) size = self.sizes[idx] ptr = self.data_offsets[idx] if length is None: length = size - offset ptr += offset a = np.empty(length, dtype=self.dtype) self.data_file.seek(ptr * self.element_size) self.data_file.readinto(a) return a def __len__(self): return self._len def num_tokens(self, index): return self.sizes[index] def size(self, index): return self.sizes[index] @staticmethod def exists(path): return os.path.exists(index_file_path(path)) and os.path.exists(data_file_path(path)) @property def supports_prefetch(self): return False # avoid prefetching to save memory @property def doc_idx(self): return self._doc_idx def get_doc_idx(self): return self._doc_idx def set_doc_idx(self, doc_idx_): self._doc_idx = doc_idx_ class IndexedDatasetBuilder(object): element_sizes = { np.uint8: 1, np.int8: 1, np.int16: 2, np.uint16: 2, np.int32: 4, np.int64: 8, np.float32: 4, np.float64: 8, } def __init__(self, out_file, dtype=np.int32): self.out_file = open(out_file, "wb") self.dtype = dtype self.data_offsets = [0] self.dim_offsets = [0] self.sizes = [] self.element_size = self.element_sizes[self.dtype] self.doc_idx = [0] def add_item(self, tensor): tensor = np.array(tensor, dtype=self.dtype) bytes = self.out_file.write(tensor) self.data_offsets.append(self.data_offsets[-1] + bytes / self.element_size) for s in tensor.shape: self.sizes.append(s) self.dim_offsets.append(self.dim_offsets[-1] + len(tensor.shape)) del bytes def end_document(self): self.doc_idx.append(len(self.sizes)) def merge_file_(self, another_file): index = IndexedDataset(another_file) assert index.dtype == self.dtype doc_offset = len(self.sizes) begin = self.data_offsets[-1] for data_offset in index.data_offsets[1:]: self.data_offsets.append(begin + data_offset) self.sizes.extend(index.sizes) begin = self.dim_offsets[-1] for dim_offset in index.dim_offsets[1:]: self.dim_offsets.append(begin + dim_offset) self.doc_idx.extend((doc_offset + index.doc_idx)[1:]) with open(data_file_path(another_file), "rb") as f: while True: data = f.read(1024) if data: self.out_file.write(data) else: break def finalize(self, index_file): self.out_file.close() index = open(index_file, "wb") index.write(b"TNTIDX\x00\x00") index.write(struct.pack(" 1 and not add_sequence_len: self._sizes.append(tensor.size) add_sequence_len = True self._data_file_dict[key].write(tensor.tobytes(order="C")) def end_document(self): self._doc_idx.append(len(self._sizes)) def finalize(self, index_file): for key, filename in self._data_file_dict.items(): filename.close() with SFTMMapIndexedDataset.Index.writer(index_file, self._dtype) as index: index.write(self._sizes, self._doc_idx) class MMapIndexedDatasetBuilder(object): def __init__(self, out_file, dtype, loss_mask_file=None): self._data_file = open(out_file, "wb") self._loss_mask_file = None if loss_mask_file is not None: self._loss_mask_file = open(loss_mask_file, "wb") self._dtype = dtype self._sizes = [] self._doc_idx = [0] def flush_loss_mask_item(self, loss_mask_lst): for loss_mask in loss_mask_lst: tensor = np.array(loss_mask, dtype=np.uint8) self._loss_mask_file.write(tensor.tobytes(order="C")) def add_item(self, tensor): tensor = np.array(tensor, dtype=self._dtype) self._data_file.write(tensor.tobytes(order="C")) self._sizes.append(tensor.size) def add_doc(self, tensor, sizes): np_array = np.array(tensor, dtype=self._dtype) self._data_file.write(np_array.tobytes(order="C")) self._sizes.extend(sizes) self._doc_idx.append(len(self._sizes)) def end_document(self): self._doc_idx.append(len(self._sizes)) def merge_file_(self, another_file): # Concatenate index index = MMapIndexedDataset.Index(index_file_path(another_file)) assert index.dtype == self._dtype offset = len(self._sizes) self._sizes.extend(index.sizes) self._doc_idx.extend((offset + index.doc_idx)[1:]) # Concatenate data with open(data_file_path(another_file), "rb") as f: shutil.copyfileobj(f, self._data_file) def finalize(self, index_file): self._data_file.close() with MMapIndexedDataset.Index.writer(index_file, self._dtype) as index: index.write(self._sizes, self._doc_idx) print("Total sentences num: %d" % len(self._sizes)) print("Total documents num: %d" % (len(self._doc_idx) - 1)) print("Total tokens num: %d" % sum(self._sizes)) print("Average tokens per sentence: %.2f" % (sum(self._sizes) / len(self._sizes))) print("Average tokens per document: %.2f" % (sum(self._sizes) / (len(self._doc_idx) - 1))) def get_indexed_dataset_(data_prefix, data_impl, skip_warmup): print_rank_0(" > building dataset index ...") start_time = time.time() indexed_dataset = make_dataset(data_prefix, data_impl, skip_warmup) assert indexed_dataset.sizes.shape[0] == indexed_dataset.doc_idx[-1] print_rank_0(" > finished creating indexed dataset in {:4f} " "seconds".format(time.time() - start_time)) print_rank_0(" > indexed dataset stats:") print_rank_0(" number of documents: {}".format(indexed_dataset.doc_idx.shape[0] - 1)) print_rank_0(" number of sentences: {}".format(indexed_dataset.sizes.shape[0])) return indexed_dataset class CompatibleIndexedDataset(paddle.io.Dataset): def __init__(self, path): super().__init__() self._path = path # All document ids, extend as 1-D array. self._token_ids = np.load(path + "_ids.npy", mmap_mode="r", allow_pickle=True) process_data = np.load(path + "_idx.npz") self._sizes = process_data["lens"] self._pointers = np.empty(len(self._sizes) + 1, dtype=np.int64) self._pointers[0] = 0 np.cumsum(self._sizes, out=self._pointers[1:]) self._doc_idx = process_data["docs"] def __getstate__(self): return self._path def __len__(self): return len(self._sizes) # @lru_cache(maxsize=8) def __getitem__(self, idx): if isinstance(idx, int): size = self._sizes[idx] ptr = self._pointers[idx] np_array = self._token_ids[ptr : ptr + size] return np_array elif isinstance(idx, slice): start, stop, step = idx.indices(len(self)) if step != 1: raise ValueError("Slices into indexed_dataset must be contiguous") ptr = self._pointers[start] sizes = self._sizes[idx] offsets = list(accumulate(sizes)) total_size = sum(sizes) np_array = self._token_ids[ptr : ptr + total_size] sents = np.split(np_array, offsets[:-1]) return sents def get(self, idx, offset=0, length=None): """Retrieves a single item from the dataset with the option to only return a portion of the item. get(idx) is the same as [idx] but get() does not support slicing. """ size = self._sizes[idx] ptr = self._pointers[idx] if length is None: length = size - offset ptr += offset np_array = self._token_ids[ptr : ptr + length] return np_array, None @property def sizes(self): return self._sizes @property def doc_idx(self): return self._doc_idx def get_doc_idx(self): return self._doc_idx def set_doc_idx(self, doc_idx_): self._doc_idx = doc_idx_ @staticmethod def exists(path): return os.path.isfile(path + "_ids.npy") and os.path.isfile(path + "_idx.npz")