uv-scripts/data-processing
Data processing Part of uv-scripts — self-contained UV scripts you run on Hugging Face Jobs in one command. General data processing recipes: convert, clean and prepare data files. Script What it does optimize-parquet.py Converts CSV, JSON and Parquet files uploaded to a bucket into optimized Parquet, triggered by a bucket webhook optimize-parquet.py: optimized Parquet from bucket uploads Upload a CSV, JSON or Parquet file to a Storage Bucket and… See the full description on the dataset page: https://huggingface.co/datasets/uv-scripts/data-processing.
051
1# /// script2# requires-python = ">=3.10"3# dependencies = ["datasets>=5"]4# ///5"""Convert CSV, JSON and Parquet files added to a bucket into optimized Parquet6in a second bucket (OUTPUT_BUCKET).7 8Meant to run as a Job triggered by a bucket webhook: the Job receives the list of9changed files in WEBHOOK_PAYLOAD. `datasets` writes optimized Parquet by default10(content-defined chunking, page index, row groups of at most 100MB).11 12Setup (once):13 14 # 1. A base Job for the webhook to re-run. With no payload, this first run exits.15 # Use `hf jobs run ... uv run <url>`, not `hf jobs uv run <url>`: the latter uploads16 # the script as a volume, and webhook runs don't keep volumes.17 hf jobs run --flavor cpu-upgrade --timeout 2h -e OUTPUT_BUCKET=<user>/<output-bucket> \\18 ghcr.io/astral-sh/uv:python3.12-bookworm \\19 uv run https://huggingface.co/datasets/uv-scripts/data-processing/raw/main/optimize-parquet.py20 21 # 2. A webhook on the input bucket that re-runs that Job on every change.22 # --secrets HF_TOKEN gives every triggered run a token for the two buckets;23 # a webhook run does not inherit the base Job's own secrets.24 hf webhooks create --job-id <job id from step 1> \\25 --watch bucket:<user>/<input-bucket> --domain repo --secrets HF_TOKEN26 27Then upload files to the input bucket, e.g.28`hf buckets cp data.csv hf://buckets/<user>/<input-bucket>/data.csv`,29and the output appears at `<output-bucket>/data.csv/data/train-00000-of-00001.parquet`.30"""31 32import json33import os34import shutil35import tempfile36from pathlib import PurePosixPath37 38from datasets import load_dataset39 40BUILDERS = {".csv": "csv", ".json": "json", ".jsonl": "json", ".parquet": "parquet"}41 42event = json.loads(os.environ.get("WEBHOOK_PAYLOAD", "{}"))43input_bucket = os.environ.get("WEBHOOK_REPO_ID")44output_bucket = os.environ["OUTPUT_BUCKET"]45# Writing to the watched bucket would trigger this Job again for its own output.46if output_bucket == input_bucket:47 raise SystemExit("OUTPUT_BUCKET must be different from the watched bucket")48 49# A full load needs disk for the download, the Arrow cache and the output.50# Larger files are streamed instead.51free_disk = shutil.disk_usage(tempfile.gettempdir()).free52stream_above = int(os.environ.get("STREAM_ABOVE_BYTES", free_disk // 3))53 54for changed_file in event.get("updatedFiles", []):55 path = PurePosixPath(changed_file["path"])56 if changed_file["action"] != "add":57 continue58 if path.suffix not in BUILDERS:59 print(f"Skipping {path}: unsupported file type")60 continue61 62 streaming = changed_file["size"] > stream_above63 mode = "streaming" if streaming else "full load"64 print(f"{path} ({changed_file['size']:,} bytes): {mode}")65 dataset = load_dataset(66 BUILDERS[path.suffix],67 data_files=f"hf://buckets/{input_bucket}/{path}",68 split="train",69 streaming=streaming,70 )71 # a/b.csv -> <output bucket>/a/b.csv/data/train-*.parquet (keeps b.csv and b.jsonl apart)72 # Tabular files have no image/audio files to embed. Setting this also avoids a73 # crash when pushing a streamed CSV/JSON dataset (its features are not known yet).74 dataset.push_to_hub(f"buckets/{output_bucket}/{path}", embed_external_files=False)75 print(f"Wrote buckets/{output_bucket}/{path}")76 