
Megatron-LM 數據集管線全解析從 IndexedDataset 二進制格式到 GPTDataset 三索引機制與快速 DataLoader 初始化【免費下載鏈接】Megatron-LMOngoing research training transformer models at scale項目地址: https://gitcode.com/GitHub_Trending/me/Megatron-LM本篇技術文章圍繞 Megatron Core 的datasets包展開系統講解其數據管線Data Pipeline的三層架構底層IndexedDataset/IndexedDatasetBuilder的二進制存取格式、中上層由BlendedMegatronDatasetConfig與BlendedMegatronDatasetBuilder驅動的分布式感知的 DataLoader 構建流程、GPTDataset的文檔/樣本/洗牌三索引查表機制以及離線緩存預生成tools/prepare_cache.py、Packing 調度器和三個加速 DataLoader 初始化的配置開關。讀完后你可以獨立完成數據集預處理、理解訓練啟動時索引緩存的生成與命中邏輯并針對大規模數據混合場景配置快速加載路徑。一、整體架構三層數據接口與構建類Megatron Core 的數據管線是分層設計的核心類關系如下源碼位于 megatron/core/datasets/層級類作用底層IndexedDataset/IndexedDatasetBuilder最低層數據接口讀寫.bin.idx二進制文件配置BlendedMegatronDatasetConfig參數化 Builder 與各級數據集可按訓練/推理體制擴展如GPTDatasetConfig構建BlendedMegatronDatasetBuilder構建最高層數據接口是分布式感知的構建入口中層MegatronDataset抽象基類建立在IndexedDataset之上的高階抽象不同任務有不同擴展如GPTDataset頂層BlendedDataset建立在多個MegatronDataset之上的混合數據集僅在多個數據分布共同貢獻同一個 split 時才需要需要注意一個重要的分布式約定所有 rank 都必須嘗試通過BlendedMegatronDatasetBuilder構建數據集否則程序會掛起哪些 rank 真正執行構建邏輯則由BlendedMegatronDatasetConfig控制。這一約定在源碼中體現為 blended_megatron_dataset_builder.py 中build_generic_dataset的 rank 0 先構建 →torch.distributed.barrier()→ 其余 rank 再構建緩存命中 模式。二、數據預處理IndexedDatasetBuilder 與 IndexedDataset數據預處理圍繞兩個類展開IndexedDatasetBuilderIndexedDataset官方文檔指出目前端到端的數據預處理實現留給用戶完成詳見類文檔字符串實際入口可參考 tools/preprocess_data.py。2.1 IndexedDatasetBuilder構建與合并數據集IndexedDatasetBuilder位于 indexed_dataset.py約 930 行起能夠構建并合并IndexedDataset實例。其核心方法包括add_item(tensor, mode0)向數據集追加單條序列并記錄其長度add_document(tensor, lengths, modesNone)追加整篇文檔lengths給出文檔內各序列的長度同時更新document_indicesadd_documents(documents, modesNone, eod_tokenNone, chunk_size1_000_000)以批量方式寫入 jagged 數組形式的文檔列表依賴awkward庫可自動在文檔邊界插入 EOD token分塊寫入以控制內存add_index(path_prefix)把另一個已存在的IndexedDataset整體合并進來拼接索引與數據finalize(idx_path)關閉數據文件并寫入.idx索引文件。構造參數為bin_path數據文件路徑、dtype默認numpy.int32與multimodal是否多模態決定是否記錄每條序列的 mode。2.2 IndexedDataset.bin 與 .idx 的二進制布局IndexedDataset是 Megatron Core 中最低層的數據接口。一個實例引用兩個二進制文件數據文件.bin保存文檔/序列數據索引文件.idx保存文檔/序列元數據。索引文件的內容分兩部分。先存數據集級元數據索引文件頭為向后兼容保留索引版本號為向后兼容保留一個數字編碼對應寫入數據文件所用的數據類型數據集中的序列數量數據集中的文檔數量。隨后存文檔級與序列級元數據按順序存每條序列的元素個數int32按順序存每條序列的字節偏移指針int64按順序存每個文檔對應的連續序列索引區間[...)int64按順序存每條序列的 mode僅多模態情形int8。這與源碼中_IndexWriter.write()的落盤順序完全一致header version dtype code sequence_count document_count sequence_lengths sequence_pointers document_indices sequence_modes其中_INDEX_HEADER bMMIDIDX\x00\x00版本號為小端8字節整型 1。數據類型編碼由DType枚舉定義int32對應編碼 4、uint16對應編碼 8 等。在讀取側_IndexReader會對.idx做numpy.memmap然后按偏移量依次還原sequence_lengthsint32、sequence_pointersint64、document_indicesint64多模態時再還原sequence_modesint8。.bin數據文件有四種讀取器實現由IndexedDataset.__init__根據參數選擇_MMapBinReadermmapTrue時內存映射讀取默認本地文件場景_FileBinReadermmapFalse時用文件指針seek readinto內置指數退避重試默認 3 次重試起始睡眠 10s 翻倍_S3BinReaderS3 對象存儲場景以bin_chunk_nbytes為塊大小維護內存緩存按塊范圍Range拉取字節_MultiStorageClientBinReader基于 Multi-Storage Client 的范圍讀。此外還暴露get(idx, offset, length)方法支持只取序列一部分的讀取這正是上層GPTDataset拼接跨文檔樣本時依賴的關鍵能力。IndexedDataset還支持fast_cache_load跳過文件存在性斷言與sequences_per_dataset直接用計數信息初始化索引免開.idx文件配合--per-dataset-sequences-path使用。三、數據加載之構建BlendedMegatronDatasetConfig 與 BuilderDataLoader 的構建是一個分布式感知的過程圍繞五個要素展開BlendedMegatronDatasetConfigBlendedMegatronDatasetBuilderIndexedDatasetMegatronDatasetBlendedDataset3.1 BlendedMegatronDatasetConfig可擴展該類blended_megatron_dataset_config.py參數化BlendedMegatronDatasetBuilder進而參數化MegatronDataset與BlendedDataset。不同訓練/推理體制需要不同的擴展例如 GPT 訓練使用的GPTDatasetConfig。配置類中的關鍵字段包括random_seed、sequence_length隨機種子與序列長度必填blend[[prefix1, prefix2], [0.3, 0.7]]形式的混合定義權重為 None 時按底層數據集長度推斷不能與blend_per_split同用blend_per_splittrain/valid/test 三個 split 各自獨立的混合定義split從單一分布抽樣時 train/valid/test 的權重字符串如99,1,0不能與blend_per_split同用path_to_cache所有可復用數據集索引的緩存目錄mmap_bin_files默認 True.bin用 mmap 還是文件指針讀取mid_level_dataset_surplus默認 0.005中層數據集構建時的樣本冗余比例頂層數據集若超采中層數據集需要調大fast_cache_load/defer_npy_index_mmap兩個快速加載開關見第七節都要求path_to_cache非空num_dataset_builder_threads構建數據集的線程數。__post_init__中完成了若干合法性約束fast_cache_load時禁止與blend同用應改用--per-split-data-args-path或 per-split data pathblend與blend_per_split互斥當blend非空時必須提供split兩者都為空時自動進入mockTrue的模擬數據模式并取1,1,1的均勻 split。split 字符串經parse_and_normalize_split歸一化后再由convert_split_vector_to_split_matrix轉換為 split 矩陣各 split 在非重疊區間上的書端點區間例如[0.99, 0.01, 0.0] - [(0, 0.99), (0.99, 1.0), None]。GPTDatasetConfiggpt_dataset.py在基類之上擴展了reset_position_ids、reset_attention_mask、eod_mask_loss、create_attention_mask、add_extra_token_to_sequence、drop_last_partial_validation_sequence、hybrid_context_parallel、sequences_per_dataset等字段并在__post_init__中依據 tokenizer 詞表大小自動推導token_dtype_code詞表超過 uint16 上限取 int32 編碼 4否則 uint16 編碼 8。3.2 BlendedMegatronDatasetBuilder最高層構建器BlendedMegatronDatasetBuilderblended_megatron_dataset_builder.py注意正確路徑為 megatron/core/datasets/blended_megatron_dataset_builder.py構建 Megatron Core 最高層的數據接口。其__init__接收四類參數cls要實例化的MegatronDataset子類sizes每個 split 要求的最少樣本總數可為 Noneis_built_on_rank判斷當前 rank 是否構建數據集的可調用對象必須感知 Megatron Core 并行策略全局 rank、組內 rank、virtual rank 都可能影響返回值且必須在全局 rank 0 上恰好返回 Trueconfig數據集配置對象。build()方法按配置分三種情形處理每個 splitsplit 為 None什么都不做單個貢獻數據集size非 None 時按比例抽樣子數據集size為 None 時不做超額抽樣多個貢獻數據集權重與 size 均給定 → 按權重與 size 構建中層與頂層數據集僅給權重不給 size → 報錯僅給 size → 頂層長度取size以各中層長度之和為上限兩者皆無 → 構建窮舉索引。值得注意的實現細節均可在源碼中驗證構建過程使用ThreadPoolExecutor并行構建各 prefix 的MegatronDataset_build_megatron_datasets_parallel每個任務通過contextvars.copy_context()把 OTel 追蹤上下文傳播進工作線程分布式初始化后默認 rank 0 先行構建線程數還會按 GPU 數適度放大barrier之后其他 rank 再構建——此時必然命中緩存開啟fast_cache_load時會跳過 rank 0 先行 barrier 的同步點各 rank 直接并行構建/加載這正是其提速原理見第七節指定了混合權重時每個 prefix 構建的樣本數由_get_size_per_split_per_dataset計算會乘以(1 mid_level_dataset_surplus)冗余系數保證頂層數據集有足夠樣本可抽。3.3 MegatronDataset可擴展與 BlendedDataset頂層MegatronDatasetmegatron_dataset.py是建立在IndexedDataset之上的高階抽象基類子類如GPTDataset需實現numel_low_level_dataset與build_low_level_dataset兩個靜態方法。構造時它會匯總類名、數據集路徑、樣本數、split、以及_key_config_attributes()返回的關鍵配置屬性random_seed、sequence_length、split、split_matrix、tokenizer序列化為 JSON 描述并計算 MD5 得到unique_description_hash——緩存索引文件名的唯一性就錨定在這個 hash 上因此任何關鍵參數變化都會導致緩存失效重建。BlendedDatasetblended_dataset.py建立在多個MegatronDataset之上僅當需要混合多個數據分布來貢獻某個 split 時才需要它混合比例通過配置控制。約束包括各子數據集必須同類、同 split、權重為正、數據集聚合數小于 32767 等若size非 None 則權重會被歸一化。四、數據加載之實現GPTDataset 的三索引機制GPTDataset由以下變量參數化底層IndexedDataset實例indexed_dataset、split 索引indexed_indices用于訓練/驗證/測試的連續文檔或序列索引子集、樣本數N、序列長度S、隨機種子R。它創建三個索引映射以支撐查表1文檔索引 Do_idx一維數組把i映射到文檔索引長度E * |indexed_indices|其中E是滿足E * |indexed_indices| N的最少 epoch 數。文檔索引按R洗牌。Given: N 15 indexed_indices [5, 6, 7, 8, 9] E 3 Then, for example: Do_idx [8, 8, 9, 6, 7, 5, 8, 5, 6, 6, 5, 9, 7, 7, 9]2樣本索引 Sa_idx二維數組把j映射到(i, Do_idx[i] 的偏移)對形狀[N 1, 2]。行j與j 1分別作為第j個樣本的左、右邊界。Given: S 1024 Then, for example: Sa_idx[0] (0, 0) Sa_idx[1] (0, 1024) Do_idx[0] has length greater than S Sa_idx[2] (1, 512) Do_idx[0] has length 1536 Sa_idx[3] (2, 0) Do_idx[1] has length 1536 Sa_idx[4] (5, 300) Do_idx[2:5] are shorter documents relative to Do_idx[0:2] Sa_idx[5] (6, 24) Do_idx[5] has length 13003洗牌索引 Sh_idx一維數組把k映射到j長度N按R洗牌。Given N 10 Then, for example: Sh_idx [4, 0, 2, 6, 1, 9, 5, 8, 7, 3]查詢第k個樣本的過程# 1. 用洗牌索引得到樣本索引內的索引 j j Sh_idx[k] # 2. 用樣本索引得到樣本左右邊界在文檔索引中的位置及各自起始 token 偏移 i, offset Sa_idx[j] i_next, offset_next Sa_idx[j 1] # 3. 用文檔索引從連續的文檔中取出 S 個 token sample [] sample indexed_dataset[Do_idx[i]][offset:] if i ! i_next: sample indexed_dataset[Do_idx[i 1:i_next]] sample indexed_dataset[Do_idx[i_next]][:offset_next]從源碼看gpt_dataset.py 的_query_document_sample_shuffle_indices約 388-473 行實際實現與偽代碼一致并做了兩點增強樣本跨越單個文檔時直接用dataset.get(idx, offset, length)一次取出跨越多個文檔時逐文檔get拼接。此外取出的 token 總數為S add_extra_token_to_sequence默認多取 1 個 token__getitem__中據此切分出tokens text[:-1]與labels text[1:]保證輸入與標簽都是完整長度不足時以 pad token 補齊并把 pad 位置的 loss mask 置零。索引構建的關鍵工程細節_build_document_sample_shuffle_indices約 475-701 行epoch 計算_get_num_epochs會不斷累加 epoch 直至 token 總量滿足N * S add_extra_token的需求末 epoch 分離若最后一個 epoch 的樣本數不足完整 epoch 樣本數的 80%threshold 0.80則把末 epoch 與前面 epoch 分開洗牌separate_final_epoch避免最后一個不完整的 epoch 被全局打散后與前面樣本過度混合C 加速樣本索引由 helpers.cpp 中的build_sample_idx構建當len(document_index) * 2 len(sequence_lengths)訪問密度高時會先把 mmap 的sequence_lengths復制進內存源碼注釋解釋了這樣做的兩個好處——順序預讀整個文件以及進入 C 時持有 GIL 提高并行度緩存三個索引與description.txt一并保存到緩存目錄文件名為{unique_description_hash}-{ClassName}-{split}-document_index.npy等fast_cache_load時跳過文件存在性檢查直接視為命中defer_npy_index_mmap時索引不在初始化時加載、延遲到首次訪問時以 mmap 方式加載此時__len__會改用純算術公式估算樣本數復用 helpers.cpp 的樣本計數邏輯。4.1 BlendedDataset 的混合索引BlendedDataset由數據集聚合D、權重W每個數據集一個與規模S參數化。它會按權重比例從各貢獻數據集抽樣直到達到目標規模每一步抽樣時從抽樣誤差sampling error最大的那個數據集抽取一個樣本。它創建兩個混合索引數據集索引 Da_idx一維數組把i映射到數據集索引長度SGiven D [d0, d1, d2] W [1/2, 1/4, 1/4] S 4 Then, for example: Da_idx [0, 1, 2, 0]數據集樣本索引 Sa_idx一維映射把i映射到數據集Da_idx[i]內的樣本索引長度SGiven Da_idx [0, 1, 2, 0] Then, for example: Sa_idx [0, 0, 0, 1]查詢第k個樣本sample D[Da_idx[k]][Sa_idx[k]]同樣為節省初始化時間各索引在單個 rank 上順序構建/緩存再由其他 rank 并行加載緩存索引錨定在BlendedDataset.__init__生成的 hash 上。源碼實現blended_dataset.py 的_build_indices中索引由 helpers.cpp 的build_blending_indicessize非 None或build_exhaustive_blending_indicessize為 None窮舉模式構建構建后還會校驗各子數據集是否被超采若超采會拋出明確提示增大mid_level_dataset_surplus的IndexError。五、離線緩存預生成tools/prepare_cache.py對于 GPT 風格訓練上述數據集緩存可以用 tools/prepare_cache.py 提前準備好而不必等訓練啟動時 rank 0 構建。該腳本復用了pretrain_gpt.py與pretrain_mamba.py的正常數據集構建路徑包括GPTDataset、BlendedDataset與BlendedMegatronDatasetBuilder。它接受常規的數據集參數支持 blend 與 per-split 數據集定義并要求--data-cache-path以便生成的緩存能被后續訓練復用。對于大型 blend 或大量文件 prefix 的場景尤其有用構建 document、sample、shuffle 索引可能耗時數分鐘期間所有 GPU 都處于空閑狀態而 rank 0 只做純 CPU 工作。如果后續訓練任務沒有指定--global-batch-size該參數用于確定數據集規模與 split應通過--prepare-cache-world-size顯式指定緩存準備時使用的 world size腳本將其直接賦給args.world_size見 prepare_cache.py 的_normalize_prepare_cache_args。明確的限制tools/prepare_cache.py不支持--mock-data、--sft、--fim-data或--step-batch-size-schedule源碼_validate_prepare_cache_args會對這些選項直接拋ValueError同時--data-cache-path為必填。腳本還會在準備階段強制關閉--dataloader-fast-cache-load與--dataloader-defer-npy-index-mmap這兩個開關的意義只在于消費已存在的緩存并在運行前打印生效的 world size、DP size、global batch size、緩存路徑與各 split 目標樣本數。六、Packing Scheduler跨 DP×CP rank 的變長序列重調度Packing 調度器把變長序列重新調度到 DP×CP 各 rank 上以提升 GPU 利用率。它圍繞以下模塊構建data_scheduledata_schedule.py 包含高層調度邏輯與入口點HybridCPDataLoaderWrapper混合上下文并行CP調度的包裝類。每次__next__調用它會(1) 從各 DP rank 拉取一批 packed 樣本(2) 在 DP 組內 all-gather 序列長度(3) 用BalancedCPScheduler來自 megatron/core/pipeline_parallel/ 的hybrid_cp_schedule調度子樣本(4) 通過 all-to-all 通信把子樣本重路由到正確的 DPxCP rank。BasePackingScheduler打包調度器的抽象基類定義了get_groups_and_subsamples()調度算法與run()完整調度流水線fetch、schedule、reroute、pack、broadcast 及 VPP 處理的接口。DpBalancedScheduler具體調度器按原始順序打包序列直到達到每個 DPxCP rank 的最大序列長度限制支持把 microbatch 數對齊到 DP size 與 VPP stage 的整數倍。wrap_data_iterator()頂層入口包裝已有的data_iterator。它創建合適的調度器、運行調度流水線、廣播元數據與新的num_microbatches返回新的數據迭代器、更新后的 microbatch 數以及 FLOPs 統計。get_batch_on_this_rank_for_sequence_packing()為當前 rank 拉取并廣播單個 packed microbatch。處理 TP/PP 廣播構造PackedSeqParams含cu_seqlens、max_seqlen、qkv_formatthd并可選地用 Transformer Engine 的thd_get_partitioned_indices在 CP rank 間劃分序列。data_schedule_utilsdata_schedule_utils.py 包含調度器使用的工具函數如broadcast_scalars、broadcast_tensor、build_packed_microbatches、reroute_samples_to_dcp_ranks等data_schedule.py頂部 import 即列明了全部依賴項。七、快速 DataLoader 初始化三個加速開關大規模訓練中DataLoader 初始化可能耗時數分鐘——因為要打開并內存映射大量文件還會顯著施壓文件系統。Megatron Core 提供了三個由配置開關控制的優化7.1 --dataloader-fast-cache-load假定數據集緩存已存在于指定的--data-cache-path中。啟用后通過移除同步點與文件檢查斷言來加速創建過程。從源碼看其具體生效點有三處配置層blended_megatron_dataset_config.py 斷言必須提供--data-cache-path且不能與--data-pathblend 形式同用應改用--per-split-data-args-path或--train-data-path/--valid-data-path/--test-data-path構建層blended_megatron_dataset_builder.py 跳過 rank 0 先行構建 barrier各 rank 直接并行構建同時跳過indexed_indices的重復計算數據層indexed_dataset.py 跳過sequence_lengths.shape[0]系列的一致性斷言。7.2 --dataloader-defer-npy-index-mmap同樣假定緩存已存在。啟用后把數據集索引.npy文件的內存映射延遲到首次訪問時進行。官方推薦與--num-workers 0搭配使用讓 DataLoader 預取下一批數據從而用后臺預取掩蓋索引 mmap 的開銷。實現上GPTDataset/BlendedDataset在_build_indices階段只記錄緩存路徑并返回 None 索引__getitem__首次調用時才safe_numpy_load(..., mmap_moder)見 gpt_dataset.py 與 blended_dataset.py__len__則改用 token 數算術公式直接推算。7.3 --per-dataset-sequences-path通過該配置指定 tools/build_sequences_per_dataset.py 生成的 JSON 文件。該腳本對 blend 中的每個文件 prefix 打開.idx讀取序列數與文檔數_IndexReader匯總為單一文件。此配置在處理數百乃至上千個文件 prefix 時尤其有用它只需要一次open操作而不是每個 prefix 一次。該 JSON 經GPTDatasetConfig.sequences_per_dataset傳入IndexedDataset后_IndexReader可跳過解析 34 字節頭部之后的完整索引讀取直接用給定的(sequence_count, document_count)初始化見 indexed_dataset.py。腳本用法示例來自其模塊 docstringpython3 tools/build_sequences_per_dataset.py --per-split-data-args-path my-training-dataset-blend.json --per-dataset-sequences-path my-training-dataset-blend-sequences-per-dataset.json八、小結一次訓練啟動中數據管線的工作順序把以上內容串起來一次 GPT 訓練啟動時數據管線的工作順序為訓練腳本通過megatron.training的參數解析生成GPTDatasetConfig含 blend、split、path_to_cache等并計算各 split 目標樣本數所有 rank 調用BlendedMegatronDatasetBuilder.build()rank 0 先構建或用--dataloader-fast-cache-load并行構建為每個 prefix 建IndexedDatasetmmap.bin解析.idx每個GPTDataset按unique_description_hash檢查緩存命中則 mmap 加載三個.npy索引或按defer_npy_index_mmap延遲加載未命中則由 rank 0 構建C 加速并寫緩存多個數據集的 split 由BlendedDataset以 最大抽樣誤差 策略生成混合索引運行時__getitem__依次經 shuffle → sample → document 三級查表拼接出定長樣本再經 DataLoader及其可選的HybridCPDataLoaderWrapperpacking 調度送入模型。相關源碼與工具入口速查主題路徑底層二進制接口megatron/core/datasets/indexed_dataset.pyC 索引構建加速megatron/core/datasets/helpers.cpp / helpers.pyGPT 數據集與三索引megatron/core/datasets/gpt_dataset.py混合數據集megatron/core/datasets/blended_dataset.py分布式構建器megatron/core/datasets/blended_megatron_dataset_builder.py配置數據類megatron/core/datasets/blended_megatron_dataset_config.py抽象基類megatron/core/datasets/megatron_dataset.py離線緩存預生成tools/prepare_cache.py每數據集元數據生成tools/build_sequences_per_dataset.pyPacking 調度器megatron/core/datasets/data_schedule.py / data_schedule_utils.py適用前提提示以上行為均基于當前倉庫版本快速加載類開關fast cache load、defer mmap、per-dataset sequences都要求緩存已預先構建且配置了--data-cache-pathtools/prepare_cache.py不支持 mock/SFT/FIM/step-batch-size-schedule 路徑離線預生成緩存時應避免這些模式。【免費下載鏈接】Megatron-LMOngoing research training transformer models at scale項目地址: https://gitcode.com/GitHub_Trending/me/Megatron-LM創作聲明:本文部分內容由AI輔助生成(AIGC),僅供參考