refactor(utils): 将 pre_tarin_process 和 sft_process 函数合并为 process_dataset

This commit is contained in:
2025-07-16 12:17:58 +08:00
parent f64a1a15ee
commit faf9f81e7c
6 changed files with 21 additions and 50 deletions
+2 -2
View File
@@ -1,9 +1,9 @@
from datasets import load_dataset from datasets import load_dataset
from utils import pre_tarin_process from utils import process_dataset
if __name__ == "__main__": if __name__ == "__main__":
dataset = load_dataset("BelleGroup/train_3.5M_CN") dataset = load_dataset("BelleGroup/train_3.5M_CN")
# pre_tarin_process( # process_dataset(
# dataset_dict=dataset, # dataset_dict=dataset,
# output_subdir="belle_sft", # output_subdir="belle_sft",
# max_chunk_size=5, # max_chunk_size=5,
+2 -2
View File
@@ -1,9 +1,9 @@
from datasets import load_dataset from datasets import load_dataset
from utils import pre_tarin_process from utils import process_dataset
if __name__ == "__main__": if __name__ == "__main__":
dataset = load_dataset("shjwudp/chinese-c4") dataset = load_dataset("shjwudp/chinese-c4")
pre_tarin_process( process_dataset(
dataset_dict=dataset, dataset_dict=dataset,
output_subdir="chinese-c4" output_subdir="chinese-c4"
) )
+2 -2
View File
@@ -1,5 +1,5 @@
from datasets import load_dataset from datasets import load_dataset
from utils import pre_tarin_process from utils import process_dataset
if __name__ == "__main__": if __name__ == "__main__":
max_chunk_num = 10 max_chunk_num = 10
@@ -10,7 +10,7 @@ if __name__ == "__main__":
data_files={"train": [f"data/0000{i}.parquet" for i in range(5)]} data_files={"train": [f"data/0000{i}.parquet" for i in range(5)]}
) )
pre_tarin_process( process_dataset(
dataset_dict=dataset, dataset_dict=dataset,
output_subdir="chinese-wiki", output_subdir="chinese-wiki",
max_chunk_num=max_chunk_num, max_chunk_num=max_chunk_num,
+2 -2
View File
@@ -1,9 +1,9 @@
from datasets import load_dataset from datasets import load_dataset
from utils import pre_tarin_process from utils import process_dataset
if __name__ == "__main__": if __name__ == "__main__":
dataset = load_dataset("HuggingFaceFW/fineweb", "sample-10BT") dataset = load_dataset("HuggingFaceFW/fineweb", "sample-10BT")
pre_tarin_process( process_dataset(
dataset_dict=dataset, dataset_dict=dataset,
output_subdir="english-fineweb", output_subdir="english-fineweb",
) )
+2 -2
View File
@@ -1,9 +1,9 @@
from datasets import load_dataset from datasets import load_dataset
from utils import pre_tarin_process from utils import process_dataset
if __name__ == "__main__": if __name__ == "__main__":
dataset = load_dataset("Blaze7451/enwiki_structured_content") dataset = load_dataset("Blaze7451/enwiki_structured_content")
pre_tarin_process( process_dataset(
dataset_dict=dataset, dataset_dict=dataset,
output_subdir="english-wiki", output_subdir="english-wiki",
max_chunk_size=5, max_chunk_size=5,
+7 -36
View File
@@ -57,13 +57,14 @@ def dump_pkl_files(
tensor = torch.cat(arrows) tensor = torch.cat(arrows)
pkl.dump(tensor, f) pkl.dump(tensor, f)
def pre_tarin_process( def process_dataset(
dataset_dict: DatasetDict, dataset_dict: DatasetDict,
output_subdir: str, output_subdir: str,
max_chunk_num: int = None, max_chunk_num: int = None,
chunk_size: int = 1000000, chunk_size: int = 1000000,
split_name: str = "train", split_name: str = "train",
column_name: str = "text", column_name: str = "text",
process_func: Callable[[dict], dict] = None,
normalization_func=comprehensive_normalization, normalization_func=comprehensive_normalization,
): ):
train_dataset = dataset_dict[split_name] train_dataset = dataset_dict[split_name]
@@ -83,43 +84,13 @@ def pre_tarin_process(
output_path = os.path.join(output_dir, f"{output_subdir}_text_chunk_{i}.jsonl") output_path = os.path.join(output_dir, f"{output_subdir}_text_chunk_{i}.jsonl")
with open(output_path, "w", encoding="utf-8") as f: with open(output_path, "w", encoding="utf-8") as f:
for example in chunk: for example in chunk:
if process_func is not None:
processed_example = process_func(example)
else:
text = example[column_name] text = example[column_name]
if normalization_func: if normalization_func:
text = normalization_func(text) text = normalization_func(text)
json_line = {column_name : text} processed_example = {column_name: text}
f.write(json.dumps(json_line, ensure_ascii=False) + "\n") f.write(json.dumps(processed_example, ensure_ascii=False) + "\n")
print(f"Saved text chunk {i} to {output_path}")
def sft_process(
dataset_dict: DatasetDict,
output_subdir: str,
max_chunk_num: int = None,
chunk_size: int = 1000000,
split_name: str = "train",
processsor: Callable[[str], str] = None,
):
train_dataset = dataset_dict[split_name]
total_samples = len(train_dataset)
num_chunks = (total_samples // chunk_size) + 1
lim_chunks = min(max_chunk_num, num_chunks) if max_chunk_num else num_chunks
script_dir = os.path.dirname(os.path.abspath(__file__))
output_dir = os.path.join(script_dir, "dataset", output_subdir)
os.makedirs(output_dir, exist_ok=True)
for i in range(lim_chunks):
start_idx = i * chunk_size
end_idx = min((i + 1) * chunk_size, total_samples)
chunk = train_dataset.select(range(start_idx, end_idx))
output_path = os.path.join(output_dir, f"{output_subdir}_text_chunk_{i}.jsonl")
with open(output_path, "w", encoding="utf-8") as f:
for example in chunk:
if processsor is not None:
example = processsor(example)
f.write(json.dumps(example, ensure_ascii=False) + "\n")
print(f"Saved text chunk {i} to {output_path}") print(f"Saved text chunk {i} to {output_path}")