应采用HuggingFace Datasets+DistributedSampler封装方案:先转换ShareGPT为Dataset对象,再经format_sharegpt统一格式、分词处理,最后用DistributedSampler切分并构建DataLoader,确保各GPU数据互斥且顺序一致。
☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 多模态理解力帮你轻松跨越从0到1的创作门槛☜☜☜

如果您在使用DeepSpeed进行分布式训练时需加载ShareGPT数据集,但发现数据无法被正确分片、出现重复样本或进程间数据不一致,则可能是由于数据集未适配DistributedSampler的索引逻辑或未处理文本长度动态性。以下是解决此问题的步骤:
一、使用HuggingFace Datasets + DistributedSampler封装
该方法通过将ShareGPT JSONL文件转换为HuggingFace Dataset对象,并结合torch.utils.data.DistributedSampler实现跨进程均匀切分,确保每个GPU仅加载互斥子集且保留原始顺序语义。
1、安装必要依赖:pip install datasets torch transformers
2、加载ShareGPT数据并构建Dataset对象:dataset = load_dataset("json", data_files={"train": "sharegpt_clean.jsonl"}, split="train")
3、定义预处理函数,统一格式化对话结构:def format_sharegpt(example): return {"text": "".join([f"### {msg['from']}: {msg['value']}" for msg in example["conversations"]])}
4、应用映射并分词:tokenized_ds = dataset.map(format_sharegpt).map(lambda x: tokenizer(x["text"], truncation=True, max_length=2048), batched=True)
5、初始化DistributedSampler:sampler = DistributedSampler(tokenized_ds, shuffle=True, drop_last=True)
6、构建DataLoader:dataloader = DataLoader(tokenized_ds, batch_size=4, sampler=sampler, num_workers=4)
二、自定义IterableDataset配合DeepSpeed的data_parallel配置
该方法适用于超大规模ShareGPT分片(如按日期/ID拆分的多个JSONL文件),避免全量加载内存,利用流式读取与rank感知路径选择实现无状态、可恢复的数据供给。
1、继承torch.utils.data.IterableDataset类,重写__iter__方法:class ShareGPTIterableDataset(IterableDataset): def __init__(self, file_list, rank, world_size): self.file_list = [f for i, f in enumerate(file_list) if i % world_size == rank]
2、在__iter__中逐行解析JSONL并yield单条样本:for file_path in self.file_list: with open(file_path) as f: for line in f: yield json.loads(line)
3、实例化数据集时传入dist.get_rank()与dist.get_world_size():ds = ShareGPTIterableDataset(glob.glob("sharegpt_*.jsonl"), dist.get_rank(), dist.get_world_size())
4、禁用sampler,直接使用DataLoader:dataloader = DataLoader(ds, batch_size=2, num_workers=2)
5、在DeepSpeed配置中显式关闭自动采样:"data_efficiency": {"enabled": false}
三、基于DeepSpeed的DataLoader Hook注入分片逻辑
该方法绕过PyTorch原生采样器,直接在DeepSpeed初始化阶段注入rank专属数据路径与偏移量,适用于已预分片且需严格控制每卡token吞吐量的场景。
1、预先将ShareGPT数据按world_size切分为独立文件:split -l 50000 sharegpt_full.jsonl sharegpt_part_
2、在init_process_group后获取当前rank对应文件:part_file = f"sharegpt_part_{dist.get_rank():02d}"
3、使用datasets.load_dataset加载该分片:local_ds = load_dataset("json", data_files=part_file, split="train")
4、调用deepspeed.initialize时传入自定义dataloader:model_engine, optimizer, _, _ = deepspeed.initialize(model=model, training_data=local_ds, ...)
5、确保DeepSpeed配置中未启用"partition_activations"或"stage3_gather_16bit_weights_on_model_save"等干扰数据流的选项:{"zero_optimization": {"stage": 2}}

















