Add Bagz format support - #4496

Open
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support
Open

Add Bagz format support#4496
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support

Conversation

@Marlon666

@Marlon666Marlon666 commented Jul 15, 2026

Copy link
Copy Markdown
Contributor

Description

This Pull Request introduces end-to-end support for the Gemini-backed Bagz (.bagz) dataset format within MaxText's PyGrain input pipeline.

Context & Problem Solved:
Previously, MaxText users relying on the ultra-efficient Bagz format for large-scale pre-training faced challenges when bridging cloud storage with PyGrain workers. This PR enables a highly optimized, POSIX-compliant file ingestion path designed explicitly for local filesystem and naturally GCS Fuse mounts. By avoiding complex REST/RPC shims, this approach achieves robust, multiprocess-safe read scalability out-of-the-box.

Scope of Changes:

  • Input Pipeline: Added the BagzDataSource wrapper in input_pipeline_utils.py for parallel and resilient Bagz shard reader instantiation.
  • Grain Integration: Registered bagz as a first-class supported file type in grain_data_processing.py, enabling native ParseFeatures and NormalizeFeatures handling.
  • Configurations: Expanded base.yml and types.py to seamlessly accept "bagz" as a valid grain_file_type.
  • Tooling: Created download_hf_dataset_as_bagz.py with multi-worker parallelism, HuggingFace streaming, and robust checkpoint recovery.

FIXES: (If there is a Github Issue #, paste it here, otherwise leave blank)

Tests

The implementation was successfully validated end-to-end on both Data Generation and Model Training levels:

  1. Parallel Shard Generation: Successfully streamed and converted the Salesforce/wikitext dataset (wikitext-103-raw-v1) into 38 contiguous .bagz file shards using multi-worker execution.
  2. End-to-End Model Training: Executed a 3-step offline training loop on a CPU-only environment using the generated local Bagz dataset via DECOUPLE_GCLOUD=TRUE (Log URL attached below).

Commands used for local testing:

# Data Generation
python tools/data_generation/download_hf_dataset_as_bagz.py \
--dataset Salesforce/wikitext \
--config '{"name": "wikitext-103-raw-v1"}' \
--output ~/datasets/wikitext_bagz \
--file-size-mb 10 \
--workers 2
# MaxText CPU Smoke Testexport DECOUPLE_GCLOUD=TRUE
python -m maxtext.trainers.pre_train.train \
src/maxtext/configs/base.yml \
run_name="bagz_cpu_smoke_test" \
steps=3 \
dataset_type="grain" \
grain_file_type="bagz" \
grain_train_files="$HOME/datasets/wikitext_bagz/*.bagz" \
grain_eval_files="$HOME/datasets/wikitext_bagz/*.bagz" \
tokenize_eval_data=True \
per_device_batch_size=1 \
max_target_length=64 \
enable_checkpointing=false \
override_model_config=true \
base_num_decoder_layers=1 \
base_emb_dim=32 \
base_mlp_dim=32 \
head_dim=8 \
base_num_kv_heads=1 \
base_num_query_heads=1
Test Logs URL: https://paste.googleplex.com/6387838210932736

Comment threadsrc/maxtext/configs/types.py
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
"""A fully picklable and fork-safe RandomAccessDataSource for Bagz files."""
def __init__(self, paths: list[str] | str):
if isinstance(paths, (list, tuple)):
self._path = ",".join(sorted(paths))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why we have to sort the paths here?

@aireenmeiaireenmei self-assigned this Jul 25, 2026

@aireenmeiaireenmei left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding the feature! Could you also add unit tests in grain_data_processing_test.py? You will need a class setup class similar to _GrainArrayRecordSetup and a test class similar to GrainArrayRecordProcessingTest

Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
# limitations under the License.

"""
Download a HuggingFace dataset via streaming and save as Bagz files.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add some example commands in the docstring?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, I've added detailed example commands in the docstring covering basic local conversion, GCSFUSE-mounted bucket output with custom prefix/file size, and authenticated access for gated datasets

Implements end-to-end support for the Bagz (.bagz) dataset format inside the
MaxText training and data processing pipeline. Enables highly efficient,
POSIX-compliant file loading optimized specifically for GCS FUSE and local
filesystem environments.
Scope of Changes:
- Input Pipeline: Integrated upstream `grain.BagzDataSource` for direct, multiprocess-safe
reading of Bagz file shards via PyGrain.
- Data Processing: Registered `bagz` as a first-class supported file type across
PyGrain data loaders, feature normalizers, and tokenizers.
- Tooling: Created `download_hf_dataset_as_bagz.py` with multi-worker support,
HuggingFace streaming, and resilient checkpointing for GCSFUSE mounts.
Testing:
- Unit Tests: Added comprehensive unit test suite in `tests/unit/bagz_data_processing_test.py`
covering RandomAccessDataSource compatibility, MapDataset, ElasticIterator (single and multi-worker),
and end-to-end `get_datasets` integration.
- End-to-End Local Smoke Test: Executed a 3-step offline training loop on a
local CPU environment using the newly generated Bagz dataset via `DECOUPLE_GCLOUD=TRUE`.
- Parallel Data Generation: Successfully converted and streamed HF datasets into
multi-shard Bagz files (38 shards, Salesforce/wikitext).
@Marlon666
Marlon666force-pushed the feat/bagz-gcsfuse-support branch from 85206e8 to 7abdb4dCompareAugust 4, 2026 21:56
…ring
Add practical example commands to the docstring of `download_hf_dataset_as_bagz.py`:
- Basic local directory conversion (Salesforce/wikitext)
- Write-optimized GCSFUSE mounted bucket output with custom shard prefix and file size
- Private/gated HuggingFace dataset download using auth token
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Marlon666@aireenmei@lepan-google
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Add Bagz format support - #4496

Open
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support
Open

Add Bagz format support#4496
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support

Conversation

@Marlon666

@Marlon666Marlon666 commented Jul 15, 2026

Copy link
Copy Markdown
Contributor

Description

This Pull Request introduces end-to-end support for the Gemini-backed Bagz (.bagz) dataset format within MaxText's PyGrain input pipeline.

Context & Problem Solved:
Previously, MaxText users relying on the ultra-efficient Bagz format for large-scale pre-training faced challenges when bridging cloud storage with PyGrain workers. This PR enables a highly optimized, POSIX-compliant file ingestion path designed explicitly for local filesystem and naturally GCS Fuse mounts. By avoiding complex REST/RPC shims, this approach achieves robust, multiprocess-safe read scalability out-of-the-box.

Scope of Changes:

  • Input Pipeline: Added the BagzDataSource wrapper in input_pipeline_utils.py for parallel and resilient Bagz shard reader instantiation.
  • Grain Integration: Registered bagz as a first-class supported file type in grain_data_processing.py, enabling native ParseFeatures and NormalizeFeatures handling.
  • Configurations: Expanded base.yml and types.py to seamlessly accept "bagz" as a valid grain_file_type.
  • Tooling: Created download_hf_dataset_as_bagz.py with multi-worker parallelism, HuggingFace streaming, and robust checkpoint recovery.

FIXES: (If there is a Github Issue #, paste it here, otherwise leave blank)

Tests

The implementation was successfully validated end-to-end on both Data Generation and Model Training levels:

  1. Parallel Shard Generation: Successfully streamed and converted the Salesforce/wikitext dataset (wikitext-103-raw-v1) into 38 contiguous .bagz file shards using multi-worker execution.
  2. End-to-End Model Training: Executed a 3-step offline training loop on a CPU-only environment using the generated local Bagz dataset via DECOUPLE_GCLOUD=TRUE (Log URL attached below).

Commands used for local testing:

# Data Generation
python tools/data_generation/download_hf_dataset_as_bagz.py \
--dataset Salesforce/wikitext \
--config '{"name": "wikitext-103-raw-v1"}' \
--output ~/datasets/wikitext_bagz \
--file-size-mb 10 \
--workers 2
# MaxText CPU Smoke Testexport DECOUPLE_GCLOUD=TRUE
python -m maxtext.trainers.pre_train.train \
src/maxtext/configs/base.yml \
run_name="bagz_cpu_smoke_test" \
steps=3 \
dataset_type="grain" \
grain_file_type="bagz" \
grain_train_files="$HOME/datasets/wikitext_bagz/*.bagz" \
grain_eval_files="$HOME/datasets/wikitext_bagz/*.bagz" \
tokenize_eval_data=True \
per_device_batch_size=1 \
max_target_length=64 \
enable_checkpointing=false \
override_model_config=true \
base_num_decoder_layers=1 \
base_emb_dim=32 \
base_mlp_dim=32 \
head_dim=8 \
base_num_kv_heads=1 \
base_num_query_heads=1
Test Logs URL: https://paste.googleplex.com/6387838210932736

Comment threadsrc/maxtext/configs/types.py
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
"""A fully picklable and fork-safe RandomAccessDataSource for Bagz files."""
def __init__(self, paths: list[str] | str):
if isinstance(paths, (list, tuple)):
self._path = ",".join(sorted(paths))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why we have to sort the paths here?

@aireenmeiaireenmei self-assigned this Jul 25, 2026

@aireenmeiaireenmei left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding the feature! Could you also add unit tests in grain_data_processing_test.py? You will need a class setup class similar to _GrainArrayRecordSetup and a test class similar to GrainArrayRecordProcessingTest

Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
# limitations under the License.

"""
Download a HuggingFace dataset via streaming and save as Bagz files.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add some example commands in the docstring?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, I've added detailed example commands in the docstring covering basic local conversion, GCSFUSE-mounted bucket output with custom prefix/file size, and authenticated access for gated datasets

Implements end-to-end support for the Bagz (.bagz) dataset format inside the
MaxText training and data processing pipeline. Enables highly efficient,
POSIX-compliant file loading optimized specifically for GCS FUSE and local
filesystem environments.
Scope of Changes:
- Input Pipeline: Integrated upstream `grain.BagzDataSource` for direct, multiprocess-safe
reading of Bagz file shards via PyGrain.
- Data Processing: Registered `bagz` as a first-class supported file type across
PyGrain data loaders, feature normalizers, and tokenizers.
- Tooling: Created `download_hf_dataset_as_bagz.py` with multi-worker support,
HuggingFace streaming, and resilient checkpointing for GCSFUSE mounts.
Testing:
- Unit Tests: Added comprehensive unit test suite in `tests/unit/bagz_data_processing_test.py`
covering RandomAccessDataSource compatibility, MapDataset, ElasticIterator (single and multi-worker),
and end-to-end `get_datasets` integration.
- End-to-End Local Smoke Test: Executed a 3-step offline training loop on a
local CPU environment using the newly generated Bagz dataset via `DECOUPLE_GCLOUD=TRUE`.
- Parallel Data Generation: Successfully converted and streamed HF datasets into
multi-shard Bagz files (38 shards, Salesforce/wikitext).
@Marlon666
Marlon666force-pushed the feat/bagz-gcsfuse-support branch from 85206e8 to 7abdb4dCompareAugust 4, 2026 21:56
…ring
Add practical example commands to the docstring of `download_hf_dataset_as_bagz.py`:
- Basic local directory conversion (Salesforce/wikitext)
- Write-optimized GCSFUSE mounted bucket output with custom shard prefix and file size
- Private/gated HuggingFace dataset download using auth token
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Marlon666@aireenmei@lepan-google
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Add Bagz format support - #4496

Open
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support
Open

Add Bagz format support#4496
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support

Conversation

@Marlon666

@Marlon666Marlon666 commented Jul 15, 2026

Copy link
Copy Markdown
Contributor

Description

This Pull Request introduces end-to-end support for the Gemini-backed Bagz (.bagz) dataset format within MaxText's PyGrain input pipeline.

Context & Problem Solved:
Previously, MaxText users relying on the ultra-efficient Bagz format for large-scale pre-training faced challenges when bridging cloud storage with PyGrain workers. This PR enables a highly optimized, POSIX-compliant file ingestion path designed explicitly for local filesystem and naturally GCS Fuse mounts. By avoiding complex REST/RPC shims, this approach achieves robust, multiprocess-safe read scalability out-of-the-box.

Scope of Changes:

  • Input Pipeline: Added the BagzDataSource wrapper in input_pipeline_utils.py for parallel and resilient Bagz shard reader instantiation.
  • Grain Integration: Registered bagz as a first-class supported file type in grain_data_processing.py, enabling native ParseFeatures and NormalizeFeatures handling.
  • Configurations: Expanded base.yml and types.py to seamlessly accept "bagz" as a valid grain_file_type.
  • Tooling: Created download_hf_dataset_as_bagz.py with multi-worker parallelism, HuggingFace streaming, and robust checkpoint recovery.

FIXES: (If there is a Github Issue #, paste it here, otherwise leave blank)

Tests

The implementation was successfully validated end-to-end on both Data Generation and Model Training levels:

  1. Parallel Shard Generation: Successfully streamed and converted the Salesforce/wikitext dataset (wikitext-103-raw-v1) into 38 contiguous .bagz file shards using multi-worker execution.
  2. End-to-End Model Training: Executed a 3-step offline training loop on a CPU-only environment using the generated local Bagz dataset via DECOUPLE_GCLOUD=TRUE (Log URL attached below).

Commands used for local testing:

# Data Generation
python tools/data_generation/download_hf_dataset_as_bagz.py \
--dataset Salesforce/wikitext \
--config '{"name": "wikitext-103-raw-v1"}' \
--output ~/datasets/wikitext_bagz \
--file-size-mb 10 \
--workers 2
# MaxText CPU Smoke Testexport DECOUPLE_GCLOUD=TRUE
python -m maxtext.trainers.pre_train.train \
src/maxtext/configs/base.yml \
run_name="bagz_cpu_smoke_test" \
steps=3 \
dataset_type="grain" \
grain_file_type="bagz" \
grain_train_files="$HOME/datasets/wikitext_bagz/*.bagz" \
grain_eval_files="$HOME/datasets/wikitext_bagz/*.bagz" \
tokenize_eval_data=True \
per_device_batch_size=1 \
max_target_length=64 \
enable_checkpointing=false \
override_model_config=true \
base_num_decoder_layers=1 \
base_emb_dim=32 \
base_mlp_dim=32 \
head_dim=8 \
base_num_kv_heads=1 \
base_num_query_heads=1
Test Logs URL: https://paste.googleplex.com/6387838210932736

Comment threadsrc/maxtext/configs/types.py
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
"""A fully picklable and fork-safe RandomAccessDataSource for Bagz files."""
def __init__(self, paths: list[str] | str):
if isinstance(paths, (list, tuple)):
self._path = ",".join(sorted(paths))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why we have to sort the paths here?

@aireenmeiaireenmei self-assigned this Jul 25, 2026

@aireenmeiaireenmei left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding the feature! Could you also add unit tests in grain_data_processing_test.py? You will need a class setup class similar to _GrainArrayRecordSetup and a test class similar to GrainArrayRecordProcessingTest

Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
# limitations under the License.

"""
Download a HuggingFace dataset via streaming and save as Bagz files.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add some example commands in the docstring?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, I've added detailed example commands in the docstring covering basic local conversion, GCSFUSE-mounted bucket output with custom prefix/file size, and authenticated access for gated datasets

Implements end-to-end support for the Bagz (.bagz) dataset format inside the
MaxText training and data processing pipeline. Enables highly efficient,
POSIX-compliant file loading optimized specifically for GCS FUSE and local
filesystem environments.
Scope of Changes:
- Input Pipeline: Integrated upstream `grain.BagzDataSource` for direct, multiprocess-safe
reading of Bagz file shards via PyGrain.
- Data Processing: Registered `bagz` as a first-class supported file type across
PyGrain data loaders, feature normalizers, and tokenizers.
- Tooling: Created `download_hf_dataset_as_bagz.py` with multi-worker support,
HuggingFace streaming, and resilient checkpointing for GCSFUSE mounts.
Testing:
- Unit Tests: Added comprehensive unit test suite in `tests/unit/bagz_data_processing_test.py`
covering RandomAccessDataSource compatibility, MapDataset, ElasticIterator (single and multi-worker),
and end-to-end `get_datasets` integration.
- End-to-End Local Smoke Test: Executed a 3-step offline training loop on a
local CPU environment using the newly generated Bagz dataset via `DECOUPLE_GCLOUD=TRUE`.
- Parallel Data Generation: Successfully converted and streamed HF datasets into
multi-shard Bagz files (38 shards, Salesforce/wikitext).
@Marlon666
Marlon666force-pushed the feat/bagz-gcsfuse-support branch from 85206e8 to 7abdb4dCompareAugust 4, 2026 21:56
…ring
Add practical example commands to the docstring of `download_hf_dataset_as_bagz.py`:
- Basic local directory conversion (Salesforce/wikitext)
- Write-optimized GCSFUSE mounted bucket output with custom shard prefix and file size
- Private/gated HuggingFace dataset download using auth token
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Marlon666@aireenmei@lepan-google
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Add Bagz format support - #4496

Open
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support
Open

Add Bagz format support#4496
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support

Conversation

@Marlon666

@Marlon666Marlon666 commented Jul 15, 2026

Copy link
Copy Markdown
Contributor

Description

This Pull Request introduces end-to-end support for the Gemini-backed Bagz (.bagz) dataset format within MaxText's PyGrain input pipeline.

Context & Problem Solved:
Previously, MaxText users relying on the ultra-efficient Bagz format for large-scale pre-training faced challenges when bridging cloud storage with PyGrain workers. This PR enables a highly optimized, POSIX-compliant file ingestion path designed explicitly for local filesystem and naturally GCS Fuse mounts. By avoiding complex REST/RPC shims, this approach achieves robust, multiprocess-safe read scalability out-of-the-box.

Scope of Changes:

  • Input Pipeline: Added the BagzDataSource wrapper in input_pipeline_utils.py for parallel and resilient Bagz shard reader instantiation.
  • Grain Integration: Registered bagz as a first-class supported file type in grain_data_processing.py, enabling native ParseFeatures and NormalizeFeatures handling.
  • Configurations: Expanded base.yml and types.py to seamlessly accept "bagz" as a valid grain_file_type.
  • Tooling: Created download_hf_dataset_as_bagz.py with multi-worker parallelism, HuggingFace streaming, and robust checkpoint recovery.

FIXES: (If there is a Github Issue #, paste it here, otherwise leave blank)

Tests

The implementation was successfully validated end-to-end on both Data Generation and Model Training levels:

  1. Parallel Shard Generation: Successfully streamed and converted the Salesforce/wikitext dataset (wikitext-103-raw-v1) into 38 contiguous .bagz file shards using multi-worker execution.
  2. End-to-End Model Training: Executed a 3-step offline training loop on a CPU-only environment using the generated local Bagz dataset via DECOUPLE_GCLOUD=TRUE (Log URL attached below).

Commands used for local testing:

# Data Generation
python tools/data_generation/download_hf_dataset_as_bagz.py \
--dataset Salesforce/wikitext \
--config '{"name": "wikitext-103-raw-v1"}' \
--output ~/datasets/wikitext_bagz \
--file-size-mb 10 \
--workers 2
# MaxText CPU Smoke Testexport DECOUPLE_GCLOUD=TRUE
python -m maxtext.trainers.pre_train.train \
src/maxtext/configs/base.yml \
run_name="bagz_cpu_smoke_test" \
steps=3 \
dataset_type="grain" \
grain_file_type="bagz" \
grain_train_files="$HOME/datasets/wikitext_bagz/*.bagz" \
grain_eval_files="$HOME/datasets/wikitext_bagz/*.bagz" \
tokenize_eval_data=True \
per_device_batch_size=1 \
max_target_length=64 \
enable_checkpointing=false \
override_model_config=true \
base_num_decoder_layers=1 \
base_emb_dim=32 \
base_mlp_dim=32 \
head_dim=8 \
base_num_kv_heads=1 \
base_num_query_heads=1
Test Logs URL: https://paste.googleplex.com/6387838210932736

Comment threadsrc/maxtext/configs/types.py
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
"""A fully picklable and fork-safe RandomAccessDataSource for Bagz files."""
def __init__(self, paths: list[str] | str):
if isinstance(paths, (list, tuple)):
self._path = ",".join(sorted(paths))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why we have to sort the paths here?

@aireenmeiaireenmei self-assigned this Jul 25, 2026

@aireenmeiaireenmei left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding the feature! Could you also add unit tests in grain_data_processing_test.py? You will need a class setup class similar to _GrainArrayRecordSetup and a test class similar to GrainArrayRecordProcessingTest

Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
# limitations under the License.

"""
Download a HuggingFace dataset via streaming and save as Bagz files.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add some example commands in the docstring?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, I've added detailed example commands in the docstring covering basic local conversion, GCSFUSE-mounted bucket output with custom prefix/file size, and authenticated access for gated datasets

Implements end-to-end support for the Bagz (.bagz) dataset format inside the
MaxText training and data processing pipeline. Enables highly efficient,
POSIX-compliant file loading optimized specifically for GCS FUSE and local
filesystem environments.
Scope of Changes:
- Input Pipeline: Integrated upstream `grain.BagzDataSource` for direct, multiprocess-safe
reading of Bagz file shards via PyGrain.
- Data Processing: Registered `bagz` as a first-class supported file type across
PyGrain data loaders, feature normalizers, and tokenizers.
- Tooling: Created `download_hf_dataset_as_bagz.py` with multi-worker support,
HuggingFace streaming, and resilient checkpointing for GCSFUSE mounts.
Testing:
- Unit Tests: Added comprehensive unit test suite in `tests/unit/bagz_data_processing_test.py`
covering RandomAccessDataSource compatibility, MapDataset, ElasticIterator (single and multi-worker),
and end-to-end `get_datasets` integration.
- End-to-End Local Smoke Test: Executed a 3-step offline training loop on a
local CPU environment using the newly generated Bagz dataset via `DECOUPLE_GCLOUD=TRUE`.
- Parallel Data Generation: Successfully converted and streamed HF datasets into
multi-shard Bagz files (38 shards, Salesforce/wikitext).
@Marlon666
Marlon666force-pushed the feat/bagz-gcsfuse-support branch from 85206e8 to 7abdb4dCompareAugust 4, 2026 21:56
…ring
Add practical example commands to the docstring of `download_hf_dataset_as_bagz.py`:
- Basic local directory conversion (Salesforce/wikitext)
- Write-optimized GCSFUSE mounted bucket output with custom shard prefix and file size
- Private/gated HuggingFace dataset download using auth token
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Marlon666@aireenmei@lepan-google
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Add Bagz format support - #4496

Open
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support
Open

Add Bagz format support#4496
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support

Conversation

@Marlon666

@Marlon666Marlon666 commented Jul 15, 2026

Copy link
Copy Markdown
Contributor

Description

This Pull Request introduces end-to-end support for the Gemini-backed Bagz (.bagz) dataset format within MaxText's PyGrain input pipeline.

Context & Problem Solved:
Previously, MaxText users relying on the ultra-efficient Bagz format for large-scale pre-training faced challenges when bridging cloud storage with PyGrain workers. This PR enables a highly optimized, POSIX-compliant file ingestion path designed explicitly for local filesystem and naturally GCS Fuse mounts. By avoiding complex REST/RPC shims, this approach achieves robust, multiprocess-safe read scalability out-of-the-box.

Scope of Changes:

  • Input Pipeline: Added the BagzDataSource wrapper in input_pipeline_utils.py for parallel and resilient Bagz shard reader instantiation.
  • Grain Integration: Registered bagz as a first-class supported file type in grain_data_processing.py, enabling native ParseFeatures and NormalizeFeatures handling.
  • Configurations: Expanded base.yml and types.py to seamlessly accept "bagz" as a valid grain_file_type.
  • Tooling: Created download_hf_dataset_as_bagz.py with multi-worker parallelism, HuggingFace streaming, and robust checkpoint recovery.

FIXES: (If there is a Github Issue #, paste it here, otherwise leave blank)

Tests

The implementation was successfully validated end-to-end on both Data Generation and Model Training levels:

  1. Parallel Shard Generation: Successfully streamed and converted the Salesforce/wikitext dataset (wikitext-103-raw-v1) into 38 contiguous .bagz file shards using multi-worker execution.
  2. End-to-End Model Training: Executed a 3-step offline training loop on a CPU-only environment using the generated local Bagz dataset via DECOUPLE_GCLOUD=TRUE (Log URL attached below).

Commands used for local testing:

# Data Generation
python tools/data_generation/download_hf_dataset_as_bagz.py \
--dataset Salesforce/wikitext \
--config '{"name": "wikitext-103-raw-v1"}' \
--output ~/datasets/wikitext_bagz \
--file-size-mb 10 \
--workers 2
# MaxText CPU Smoke Testexport DECOUPLE_GCLOUD=TRUE
python -m maxtext.trainers.pre_train.train \
src/maxtext/configs/base.yml \
run_name="bagz_cpu_smoke_test" \
steps=3 \
dataset_type="grain" \
grain_file_type="bagz" \
grain_train_files="$HOME/datasets/wikitext_bagz/*.bagz" \
grain_eval_files="$HOME/datasets/wikitext_bagz/*.bagz" \
tokenize_eval_data=True \
per_device_batch_size=1 \
max_target_length=64 \
enable_checkpointing=false \
override_model_config=true \
base_num_decoder_layers=1 \
base_emb_dim=32 \
base_mlp_dim=32 \
head_dim=8 \
base_num_kv_heads=1 \
base_num_query_heads=1
Test Logs URL: https://paste.googleplex.com/6387838210932736

Comment threadsrc/maxtext/configs/types.py
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
"""A fully picklable and fork-safe RandomAccessDataSource for Bagz files."""
def __init__(self, paths: list[str] | str):
if isinstance(paths, (list, tuple)):
self._path = ",".join(sorted(paths))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why we have to sort the paths here?

@aireenmeiaireenmei self-assigned this Jul 25, 2026

@aireenmeiaireenmei left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding the feature! Could you also add unit tests in grain_data_processing_test.py? You will need a class setup class similar to _GrainArrayRecordSetup and a test class similar to GrainArrayRecordProcessingTest

Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
# limitations under the License.

"""
Download a HuggingFace dataset via streaming and save as Bagz files.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add some example commands in the docstring?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, I've added detailed example commands in the docstring covering basic local conversion, GCSFUSE-mounted bucket output with custom prefix/file size, and authenticated access for gated datasets

Implements end-to-end support for the Bagz (.bagz) dataset format inside the
MaxText training and data processing pipeline. Enables highly efficient,
POSIX-compliant file loading optimized specifically for GCS FUSE and local
filesystem environments.
Scope of Changes:
- Input Pipeline: Integrated upstream `grain.BagzDataSource` for direct, multiprocess-safe
reading of Bagz file shards via PyGrain.
- Data Processing: Registered `bagz` as a first-class supported file type across
PyGrain data loaders, feature normalizers, and tokenizers.
- Tooling: Created `download_hf_dataset_as_bagz.py` with multi-worker support,
HuggingFace streaming, and resilient checkpointing for GCSFUSE mounts.
Testing:
- Unit Tests: Added comprehensive unit test suite in `tests/unit/bagz_data_processing_test.py`
covering RandomAccessDataSource compatibility, MapDataset, ElasticIterator (single and multi-worker),
and end-to-end `get_datasets` integration.
- End-to-End Local Smoke Test: Executed a 3-step offline training loop on a
local CPU environment using the newly generated Bagz dataset via `DECOUPLE_GCLOUD=TRUE`.
- Parallel Data Generation: Successfully converted and streamed HF datasets into
multi-shard Bagz files (38 shards, Salesforce/wikitext).
@Marlon666
Marlon666force-pushed the feat/bagz-gcsfuse-support branch from 85206e8 to 7abdb4dCompareAugust 4, 2026 21:56
…ring
Add practical example commands to the docstring of `download_hf_dataset_as_bagz.py`:
- Basic local directory conversion (Salesforce/wikitext)
- Write-optimized GCSFUSE mounted bucket output with custom shard prefix and file size
- Private/gated HuggingFace dataset download using auth token
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Marlon666@aireenmei@lepan-google
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Add Bagz format support - #4496

Open
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support
Open

Add Bagz format support#4496
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support

Conversation

@Marlon666

@Marlon666Marlon666 commented Jul 15, 2026

Copy link
Copy Markdown
Contributor

Description

This Pull Request introduces end-to-end support for the Gemini-backed Bagz (.bagz) dataset format within MaxText's PyGrain input pipeline.

Context & Problem Solved:
Previously, MaxText users relying on the ultra-efficient Bagz format for large-scale pre-training faced challenges when bridging cloud storage with PyGrain workers. This PR enables a highly optimized, POSIX-compliant file ingestion path designed explicitly for local filesystem and naturally GCS Fuse mounts. By avoiding complex REST/RPC shims, this approach achieves robust, multiprocess-safe read scalability out-of-the-box.

Scope of Changes:

  • Input Pipeline: Added the BagzDataSource wrapper in input_pipeline_utils.py for parallel and resilient Bagz shard reader instantiation.
  • Grain Integration: Registered bagz as a first-class supported file type in grain_data_processing.py, enabling native ParseFeatures and NormalizeFeatures handling.
  • Configurations: Expanded base.yml and types.py to seamlessly accept "bagz" as a valid grain_file_type.
  • Tooling: Created download_hf_dataset_as_bagz.py with multi-worker parallelism, HuggingFace streaming, and robust checkpoint recovery.

FIXES: (If there is a Github Issue #, paste it here, otherwise leave blank)

Tests

The implementation was successfully validated end-to-end on both Data Generation and Model Training levels:

  1. Parallel Shard Generation: Successfully streamed and converted the Salesforce/wikitext dataset (wikitext-103-raw-v1) into 38 contiguous .bagz file shards using multi-worker execution.
  2. End-to-End Model Training: Executed a 3-step offline training loop on a CPU-only environment using the generated local Bagz dataset via DECOUPLE_GCLOUD=TRUE (Log URL attached below).

Commands used for local testing:

# Data Generation
python tools/data_generation/download_hf_dataset_as_bagz.py \
--dataset Salesforce/wikitext \
--config '{"name": "wikitext-103-raw-v1"}' \
--output ~/datasets/wikitext_bagz \
--file-size-mb 10 \
--workers 2
# MaxText CPU Smoke Testexport DECOUPLE_GCLOUD=TRUE
python -m maxtext.trainers.pre_train.train \
src/maxtext/configs/base.yml \
run_name="bagz_cpu_smoke_test" \
steps=3 \
dataset_type="grain" \
grain_file_type="bagz" \
grain_train_files="$HOME/datasets/wikitext_bagz/*.bagz" \
grain_eval_files="$HOME/datasets/wikitext_bagz/*.bagz" \
tokenize_eval_data=True \
per_device_batch_size=1 \
max_target_length=64 \
enable_checkpointing=false \
override_model_config=true \
base_num_decoder_layers=1 \
base_emb_dim=32 \
base_mlp_dim=32 \
head_dim=8 \
base_num_kv_heads=1 \
base_num_query_heads=1
Test Logs URL: https://paste.googleplex.com/6387838210932736

Comment threadsrc/maxtext/configs/types.py
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
"""A fully picklable and fork-safe RandomAccessDataSource for Bagz files."""
def __init__(self, paths: list[str] | str):
if isinstance(paths, (list, tuple)):
self._path = ",".join(sorted(paths))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why we have to sort the paths here?

@aireenmeiaireenmei self-assigned this Jul 25, 2026

@aireenmeiaireenmei left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding the feature! Could you also add unit tests in grain_data_processing_test.py? You will need a class setup class similar to _GrainArrayRecordSetup and a test class similar to GrainArrayRecordProcessingTest

Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
# limitations under the License.

"""
Download a HuggingFace dataset via streaming and save as Bagz files.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add some example commands in the docstring?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, I've added detailed example commands in the docstring covering basic local conversion, GCSFUSE-mounted bucket output with custom prefix/file size, and authenticated access for gated datasets

Implements end-to-end support for the Bagz (.bagz) dataset format inside the
MaxText training and data processing pipeline. Enables highly efficient,
POSIX-compliant file loading optimized specifically for GCS FUSE and local
filesystem environments.
Scope of Changes:
- Input Pipeline: Integrated upstream `grain.BagzDataSource` for direct, multiprocess-safe
reading of Bagz file shards via PyGrain.
- Data Processing: Registered `bagz` as a first-class supported file type across
PyGrain data loaders, feature normalizers, and tokenizers.
- Tooling: Created `download_hf_dataset_as_bagz.py` with multi-worker support,
HuggingFace streaming, and resilient checkpointing for GCSFUSE mounts.
Testing:
- Unit Tests: Added comprehensive unit test suite in `tests/unit/bagz_data_processing_test.py`
covering RandomAccessDataSource compatibility, MapDataset, ElasticIterator (single and multi-worker),
and end-to-end `get_datasets` integration.
- End-to-End Local Smoke Test: Executed a 3-step offline training loop on a
local CPU environment using the newly generated Bagz dataset via `DECOUPLE_GCLOUD=TRUE`.
- Parallel Data Generation: Successfully converted and streamed HF datasets into
multi-shard Bagz files (38 shards, Salesforce/wikitext).
@Marlon666
Marlon666force-pushed the feat/bagz-gcsfuse-support branch from 85206e8 to 7abdb4dCompareAugust 4, 2026 21:56
…ring
Add practical example commands to the docstring of `download_hf_dataset_as_bagz.py`:
- Basic local directory conversion (Salesforce/wikitext)
- Write-optimized GCSFUSE mounted bucket output with custom shard prefix and file size
- Private/gated HuggingFace dataset download using auth token
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Marlon666@aireenmei@lepan-google
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Add Bagz format support - #4496

Open
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support
Open

Add Bagz format support#4496
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support

Conversation

@Marlon666

@Marlon666Marlon666 commented Jul 15, 2026

Copy link
Copy Markdown
Contributor

Description

This Pull Request introduces end-to-end support for the Gemini-backed Bagz (.bagz) dataset format within MaxText's PyGrain input pipeline.

Context & Problem Solved:
Previously, MaxText users relying on the ultra-efficient Bagz format for large-scale pre-training faced challenges when bridging cloud storage with PyGrain workers. This PR enables a highly optimized, POSIX-compliant file ingestion path designed explicitly for local filesystem and naturally GCS Fuse mounts. By avoiding complex REST/RPC shims, this approach achieves robust, multiprocess-safe read scalability out-of-the-box.

Scope of Changes:

  • Input Pipeline: Added the BagzDataSource wrapper in input_pipeline_utils.py for parallel and resilient Bagz shard reader instantiation.
  • Grain Integration: Registered bagz as a first-class supported file type in grain_data_processing.py, enabling native ParseFeatures and NormalizeFeatures handling.
  • Configurations: Expanded base.yml and types.py to seamlessly accept "bagz" as a valid grain_file_type.
  • Tooling: Created download_hf_dataset_as_bagz.py with multi-worker parallelism, HuggingFace streaming, and robust checkpoint recovery.

FIXES: (If there is a Github Issue #, paste it here, otherwise leave blank)

Tests

The implementation was successfully validated end-to-end on both Data Generation and Model Training levels:

  1. Parallel Shard Generation: Successfully streamed and converted the Salesforce/wikitext dataset (wikitext-103-raw-v1) into 38 contiguous .bagz file shards using multi-worker execution.
  2. End-to-End Model Training: Executed a 3-step offline training loop on a CPU-only environment using the generated local Bagz dataset via DECOUPLE_GCLOUD=TRUE (Log URL attached below).

Commands used for local testing:

# Data Generation
python tools/data_generation/download_hf_dataset_as_bagz.py \
--dataset Salesforce/wikitext \
--config '{"name": "wikitext-103-raw-v1"}' \
--output ~/datasets/wikitext_bagz \
--file-size-mb 10 \
--workers 2
# MaxText CPU Smoke Testexport DECOUPLE_GCLOUD=TRUE
python -m maxtext.trainers.pre_train.train \
src/maxtext/configs/base.yml \
run_name="bagz_cpu_smoke_test" \
steps=3 \
dataset_type="grain" \
grain_file_type="bagz" \
grain_train_files="$HOME/datasets/wikitext_bagz/*.bagz" \
grain_eval_files="$HOME/datasets/wikitext_bagz/*.bagz" \
tokenize_eval_data=True \
per_device_batch_size=1 \
max_target_length=64 \
enable_checkpointing=false \
override_model_config=true \
base_num_decoder_layers=1 \
base_emb_dim=32 \
base_mlp_dim=32 \
head_dim=8 \
base_num_kv_heads=1 \
base_num_query_heads=1
Test Logs URL: https://paste.googleplex.com/6387838210932736

Comment threadsrc/maxtext/configs/types.py
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
"""A fully picklable and fork-safe RandomAccessDataSource for Bagz files."""
def __init__(self, paths: list[str] | str):
if isinstance(paths, (list, tuple)):
self._path = ",".join(sorted(paths))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why we have to sort the paths here?

@aireenmeiaireenmei self-assigned this Jul 25, 2026

@aireenmeiaireenmei left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding the feature! Could you also add unit tests in grain_data_processing_test.py? You will need a class setup class similar to _GrainArrayRecordSetup and a test class similar to GrainArrayRecordProcessingTest

Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
# limitations under the License.

"""
Download a HuggingFace dataset via streaming and save as Bagz files.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add some example commands in the docstring?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, I've added detailed example commands in the docstring covering basic local conversion, GCSFUSE-mounted bucket output with custom prefix/file size, and authenticated access for gated datasets

Implements end-to-end support for the Bagz (.bagz) dataset format inside the
MaxText training and data processing pipeline. Enables highly efficient,
POSIX-compliant file loading optimized specifically for GCS FUSE and local
filesystem environments.
Scope of Changes:
- Input Pipeline: Integrated upstream `grain.BagzDataSource` for direct, multiprocess-safe
reading of Bagz file shards via PyGrain.
- Data Processing: Registered `bagz` as a first-class supported file type across
PyGrain data loaders, feature normalizers, and tokenizers.
- Tooling: Created `download_hf_dataset_as_bagz.py` with multi-worker support,
HuggingFace streaming, and resilient checkpointing for GCSFUSE mounts.
Testing:
- Unit Tests: Added comprehensive unit test suite in `tests/unit/bagz_data_processing_test.py`
covering RandomAccessDataSource compatibility, MapDataset, ElasticIterator (single and multi-worker),
and end-to-end `get_datasets` integration.
- End-to-End Local Smoke Test: Executed a 3-step offline training loop on a
local CPU environment using the newly generated Bagz dataset via `DECOUPLE_GCLOUD=TRUE`.
- Parallel Data Generation: Successfully converted and streamed HF datasets into
multi-shard Bagz files (38 shards, Salesforce/wikitext).
@Marlon666
Marlon666force-pushed the feat/bagz-gcsfuse-support branch from 85206e8 to 7abdb4dCompareAugust 4, 2026 21:56
…ring
Add practical example commands to the docstring of `download_hf_dataset_as_bagz.py`:
- Basic local directory conversion (Salesforce/wikitext)
- Write-optimized GCSFUSE mounted bucket output with custom shard prefix and file size
- Private/gated HuggingFace dataset download using auth token
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Marlon666@aireenmei@lepan-google
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Add Bagz format support - #4496

Open
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support
Open

Add Bagz format support#4496
Marlon666 wants to merge 2 commits into
AI-Hypercomputer:mainfrom
Marlon666:feat/bagz-gcsfuse-support

Conversation

@Marlon666

@Marlon666Marlon666 commented Jul 15, 2026

Copy link
Copy Markdown
Contributor

Description

This Pull Request introduces end-to-end support for the Gemini-backed Bagz (.bagz) dataset format within MaxText's PyGrain input pipeline.

Context & Problem Solved:
Previously, MaxText users relying on the ultra-efficient Bagz format for large-scale pre-training faced challenges when bridging cloud storage with PyGrain workers. This PR enables a highly optimized, POSIX-compliant file ingestion path designed explicitly for local filesystem and naturally GCS Fuse mounts. By avoiding complex REST/RPC shims, this approach achieves robust, multiprocess-safe read scalability out-of-the-box.

Scope of Changes:

  • Input Pipeline: Added the BagzDataSource wrapper in input_pipeline_utils.py for parallel and resilient Bagz shard reader instantiation.
  • Grain Integration: Registered bagz as a first-class supported file type in grain_data_processing.py, enabling native ParseFeatures and NormalizeFeatures handling.
  • Configurations: Expanded base.yml and types.py to seamlessly accept "bagz" as a valid grain_file_type.
  • Tooling: Created download_hf_dataset_as_bagz.py with multi-worker parallelism, HuggingFace streaming, and robust checkpoint recovery.

FIXES: (If there is a Github Issue #, paste it here, otherwise leave blank)

Tests

The implementation was successfully validated end-to-end on both Data Generation and Model Training levels:

  1. Parallel Shard Generation: Successfully streamed and converted the Salesforce/wikitext dataset (wikitext-103-raw-v1) into 38 contiguous .bagz file shards using multi-worker execution.
  2. End-to-End Model Training: Executed a 3-step offline training loop on a CPU-only environment using the generated local Bagz dataset via DECOUPLE_GCLOUD=TRUE (Log URL attached below).

Commands used for local testing:

# Data Generation
python tools/data_generation/download_hf_dataset_as_bagz.py \
--dataset Salesforce/wikitext \
--config '{"name": "wikitext-103-raw-v1"}' \
--output ~/datasets/wikitext_bagz \
--file-size-mb 10 \
--workers 2
# MaxText CPU Smoke Testexport DECOUPLE_GCLOUD=TRUE
python -m maxtext.trainers.pre_train.train \
src/maxtext/configs/base.yml \
run_name="bagz_cpu_smoke_test" \
steps=3 \
dataset_type="grain" \
grain_file_type="bagz" \
grain_train_files="$HOME/datasets/wikitext_bagz/*.bagz" \
grain_eval_files="$HOME/datasets/wikitext_bagz/*.bagz" \
tokenize_eval_data=True \
per_device_batch_size=1 \
max_target_length=64 \
enable_checkpointing=false \
override_model_config=true \
base_num_decoder_layers=1 \
base_emb_dim=32 \
base_mlp_dim=32 \
head_dim=8 \
base_num_kv_heads=1 \
base_num_query_heads=1
Test Logs URL: https://paste.googleplex.com/6387838210932736

Comment threadsrc/maxtext/configs/types.py
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
"""A fully picklable and fork-safe RandomAccessDataSource for Bagz files."""
def __init__(self, paths: list[str] | str):
if isinstance(paths, (list, tuple)):
self._path = ",".join(sorted(paths))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why we have to sort the paths here?

@aireenmeiaireenmei self-assigned this Jul 25, 2026

@aireenmeiaireenmei left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding the feature! Could you also add unit tests in grain_data_processing_test.py? You will need a class setup class similar to _GrainArrayRecordSetup and a test class similar to GrainArrayRecordProcessingTest

Comment threadsrc/maxtext/input_pipeline/grain_data_processing.py Outdated
# limitations under the License.

"""
Download a HuggingFace dataset via streaming and save as Bagz files.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add some example commands in the docstring?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, I've added detailed example commands in the docstring covering basic local conversion, GCSFUSE-mounted bucket output with custom prefix/file size, and authenticated access for gated datasets

Implements end-to-end support for the Bagz (.bagz) dataset format inside the
MaxText training and data processing pipeline. Enables highly efficient,
POSIX-compliant file loading optimized specifically for GCS FUSE and local
filesystem environments.
Scope of Changes:
- Input Pipeline: Integrated upstream `grain.BagzDataSource` for direct, multiprocess-safe
reading of Bagz file shards via PyGrain.
- Data Processing: Registered `bagz` as a first-class supported file type across
PyGrain data loaders, feature normalizers, and tokenizers.
- Tooling: Created `download_hf_dataset_as_bagz.py` with multi-worker support,
HuggingFace streaming, and resilient checkpointing for GCSFUSE mounts.
Testing:
- Unit Tests: Added comprehensive unit test suite in `tests/unit/bagz_data_processing_test.py`
covering RandomAccessDataSource compatibility, MapDataset, ElasticIterator (single and multi-worker),
and end-to-end `get_datasets` integration.
- End-to-End Local Smoke Test: Executed a 3-step offline training loop on a
local CPU environment using the newly generated Bagz dataset via `DECOUPLE_GCLOUD=TRUE`.
- Parallel Data Generation: Successfully converted and streamed HF datasets into
multi-shard Bagz files (38 shards, Salesforce/wikitext).
@Marlon666
Marlon666force-pushed the feat/bagz-gcsfuse-support branch from 85206e8 to 7abdb4dCompareAugust 4, 2026 21:56
…ring
Add practical example commands to the docstring of `download_hf_dataset_as_bagz.py`:
- Basic local directory conversion (Salesforce/wikitext)
- Write-optimized GCSFUSE mounted bucket output with custom shard prefix and file size
- Private/gated HuggingFace dataset download using auth token
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Marlon666@aireenmei@lepan-google