
Accelerated Computing Cudf
- 3 installs
- 7 repo stars
- Updated August 2, 2026
- practicalswan/agent-skills
accelerated-computing-cudf is a Claude Code skill for ai & agent building.
About
Guides an implementer through correct, fast GPU DataFrame code using cudf.pandas or explicit cuDF for ETL, joins, groupby, and CSV/Parquet I/O. A data engineer uses it when migrating pandas workloads to CUDA GPUs or building multi-GPU DataFrame pipelines.
- Chooses cudf.pandas vs explicit cuDF based on user intent
- Tracks CUDA/driver compatibility and cuDF release 26.04 semantics
Accelerated Computing Cudf by the numbers
- 3 all-time installs (skills.sh)
- Ranked #1,661 of 2,064 Data Science & ML skills by installs in the Skillselion catalog
- Data as of Aug 4, 2026 (Skillselion catalog sync)
npx skills add https://github.com/practicalswan/agent-skills --skill accelerated-computing-cudfAdd your badge
Show developers this skill is listed on Skillselion. Paste this into your README.
| Installs | 3 |
|---|---|
| repo stars | ★ 7 |
| Last updated | August 2, 2026 |
| Repository | practicalswan/agent-skills ↗ |
How do I helps with ai & agent building tasks.?
NVIDIA-authored guidance for cuDF GPU DataFrames, pandas acceleration, dask-cuDF, and ETL joins/groupby with nullable and I/O semantics.
Who is it for?
A solo builder working on ai & agent building tasks who needs structured help with accelerated computing cudf.
Skip if: Teams with no ai & agent building needs, or anyone wanting a generic chat assistant without this specific workflow.
When should I use this skill?
When you need to helps with ai & agent building tasks., or when accelerated-computing-cudf is a claude code skill for ai & agent building.
What you get
Structured output aligned to accelerated-computing-cudf: accelerated-computing-cudf, AI & Agent Building.
Files
cuDF & dask-cuDF Implementer's Guide
Compatibility
- Release tracked by this skill: 26.04.
- Requires NVIDIA Volta or newer on CUDA 12, or Turing or newer on CUDA 13. Release 26.04 supports CUDA 12.2-12.9 with driver 535+ or CUDA 13.0-13.1 with driver 580+, and Python 3.11-3.14. cuDF sweet spot: >100K rows.
Naming
Use NVIDIA library-first wording in user-facing answers. Keep literal RAPIDS/rapidsai URLs, package names, and release metadata when citing sources.
Role
You are a cuDF expert helping an implementer work with GPU DataFrames. The user understands pandas and their data — your job is to get them to correct, fast GPU code with minimal friction. Choose the path from the user's intent: cudf.pandas for broad compatibility or minimal-change acceleration, explicit cuDF for named DataFrame migrations, hot ETL paths, and parity-sensitive work. Treat source schema, row counts, null placement, ordering, and numeric tolerances as user-visible behavior.
Critical Rules
1. Choose the right cuDF path. Use cudf.pandas for broad compatibility or minimal-change acceleration. Use explicit cuDF when the user asks to migrate DataFrame code, inspect parity, optimize a visible ETL hot path, or control unsupported operations. 2. Size gate: 100K rows minimum. Below that, GPU transfer overhead usually beats the speedup; use small data for correctness and benchmark larger working sets for performance. 3. Keep conversions at boundaries. Use .to_pandas(), .values, or .numpy() for display, plotting, CPU-only libraries, or final output boundaries. Keep intermediate ETL data on GPU. 4. Float32 is your friend. cuDF operations on float64 are slower; cast early when precision allows. 5. Validate semantics on representative slices. For null handling, joins, time series, reshape, or grouped logic, keep a small pandas reference path and compare shape, labels, null counts, ordering, and representative values before claiming parity. 6. For data > GPU memory, move to dask-cuDF with enable_cudf_spill=True. See references/dask-cudf-patterns.md.
Three Paths to GPU DataFrames
Path 1: cudf.pandas Accelerator (Compatibility / Minimal Change)
Use when the user needs a small code change, third-party pandas compatibility, or one code path that can keep running while unsupported operations fall back.
Jupyter/IPython:
%load_ext cudf.pandas
import pandas as pd # now GPU-backed; falls back silently for unsupported opsScript:
python -m cudf.pandas my_script.pyWith multiprocessing:
import cudf.pandas
cudf.pandas.install() # must come BEFORE pandas import, before Pool creation
from multiprocessing import PoolConfirm acceleration with the cudf.pandas profiler before claiming speedup. For notebook, CLI, and stats examples, read references/cudf-pandas-accelerator.md. If the profile shows the hot path running on CPU, use Path 2 for explicit cuDF control.
Path 2: Explicit cuDF API
For full control, hot-path optimization, named DataFrame migrations, and parity-sensitive operations:
import cudf
# Read data directly to GPU
df = cudf.read_parquet("data.parquet")
# Operations mirror pandas
result = df.groupby("key")["value"].sum()
merged = df.merge(lookup, on="id", how="left")
filtered = df[df["amount"] > 1000]
# String operations
df["clean"] = df["name"].str.strip().str.lower()
# To check API coverage before committing to migration:
# See references/api-patterns.md for known gaps and workaroundsKeep data on GPU end-to-end. Only call .to_pandas() at the very end for display or CPU or non-GPU handoff.
Prefer explicit cuDF for tasks involving read_csv/read_parquet, joins, groupby, reshape, nullable types, fillna/where, time buckets, rolling windows, or CPU/GPU parity checks. Add a small CPU/GPU validation path when semantics matter instead of relying on successful execution alone.
For pandas code with null handling, reshape, or time-series behavior, read references/api-patterns.md for the relevant semantic checklist before rewriting. A cudf.pandas bootstrap is enough for a minimal-change request; an implementation request should make the hot path explicit and observable.
For reshape-heavy pandas code (pivot_table, melt, stack/unstack, crosstab), keep the source schema as part of the contract: index labels, column labels or levels, fill_value, aggfunc, margins, and normalization. Use explicit cuDF where the equivalent is supported; use cudf.pandas or a narrow compatibility boundary when exact pandas reshape semantics matter more than rewriting every operation. Add a small pandas-reference parity check for shape, labels, and representative values before finalizing. See references/api-patterns.md.
Path 3: dask-cuDF (Multi-GPU / Large Data)
When dataset exceeds GPU memory. See references/dask-cudf-patterns.md for full patterns.
from dask_cuda import LocalCUDACluster
from dask.distributed import Client
import dask_cudf
cluster = LocalCUDACluster(enable_cudf_spill=True) # one worker per GPU
client = Client(cluster)
ddf = dask_cudf.read_parquet("s3://bucket/data/*.parquet")
result = ddf.groupby("key").agg({"value": "sum"}).compute()Memory Management
Enable spill before OOM happens (not after):
import cudf
cudf.set_option("spill", True) # spill to host RAM when GPU is fullRMM pool allocator (reduces cudaMalloc overhead in pipelines with many allocations):
import rmm
rmm.set_current_device_resource(rmm.mr.CudaAsyncMemoryResource())
# Must be called BEFORE any cuDF operations| GPU Free vs Dataset | Strategy |
|---|---|
| Free > 2× dataset | Single GPU cuDF |
| Free 1–2× dataset | cuDF + cudf.set_option("spill", True) |
| Dataset > GPU mem | dask-cuDF |
| Dataset > node mem | dask-cuDF + multi-node (see accelerated-computing-mpf) |
Troubleshooting
No speedup vs pandas:
- Data < 100K rows? GPU overhead dominates, so treat the run as correctness validation and measure speedup on a larger working set.
- Run
%%cudf.pandas.profile— high CPU % means many fallbacks. Identify and fix those ops. - Check
references/api-patterns.mdfor known gaps.
OOM (CUDA out of memory): 1. Enable spill: cudf.set_option("spill", True) 2. If allocator fragmentation or repeated allocation overhead is visible, use the accelerated-computing-rmm memory-resource setup guidance before GPU allocations 3. Still failing: move to dask-cuDF
AttributeError / NotImplementedError:
- Check
references/api-patterns.mdfor the specific operation - Keep that one operation on CPU at a narrow boundary and continue the supported pipeline on GPU
- Use
.to_pandas()only for the unsupported op, then.from_pandas()back
Wrong results vs pandas:
- Null/NaN handling differs: cuDF uses
<NA>(nullable) by default, pandas usesNaN. Seereferences/api-patterns.md. - Sort stability: cuDF sort is not guaranteed stable unless
stable=Trueis passed - If the difference is due to floating point differences, try casting to higher precision floats (e.g.
float64instead offloat32). If the results are still different, stop. GPU and CPU algorithms will always produce different results on floating point numbers due to the non-associativity of floating point arithmetic and that cannot be fixed.
Nullable and Fill Semantics
When the user explicitly cares about pandas nullable dtypes, fillna, where/mask, or grouped null behavior, treat parity checks as part of the implementation. See references/api-patterns.md for nullable dtype examples.
- Preserve nullable integer/string columns instead of filling them with sentinel
values unless the source code already did that.
- Keep
where/masksemantics when they encode a condition. Use broad
fillna only when the condition is exactly null-only.
- Compare with
to_pandas(nullable=True)when the pandas reference uses
nullable extension dtypes.
- Put the parity check in a reusable helper next to the GPU path, so future
changes exercise the same nullable conversion and aggregation checks.
- Validate row counts, null counts, mask truth tables, grouped aggregates, and
representative dtypes before claiming semantic parity.
Reference Files
references/cudf-pandas-accelerator.md— Profiling, fallback detection, cudf.pandas deep divereferences/api-patterns.md— Known API gaps, workarounds, semantic differencesreferences/dask-cudf-patterns.md— Multi-GPU patterns, best practices, partition tuning
External Documentation
Use WebFetch to retrieve detailed API signatures, parameter descriptions, and examples on demand.
- cuDF Documentation: https://docs.rapids.ai/api/cudf/stable/
- dask-cuDF API Reference: https://docs.rapids.ai/api/dask-cudf/stable/api/
- GitHub: https://github.com/rapidsai/cudf
- CHANGELOG: https://github.com/rapidsai/cudf/blob/main/CHANGELOG.md
Anti-Patterns
- Assuming the GPU path wins on tiny datasets: Below the size gate, transfer overhead can erase the gain.
- Bouncing between pandas and cuDF mid-pipeline: Repeated boundary conversions hide the real bottleneck and muddy parity checks.
- Claiming migration success without representative semantics checks: Null handling, ordering, and dtype drift are user-visible bugs.
Verification Protocol
Before claiming "skill applied successfully":
1. Pass/fail: The answer chooses the correct acceleration path (cudf.pandas, explicit cuDF, or dask-cuDF) for the stated workload. 2. Pass/fail: CPU/GPU conversion boundaries, dtype assumptions, and null or ordering risks are called out before claiming parity. 3. Pass/fail: Any performance claim includes a realistic dataset size and a reproducible benchmark or profiler path. 4. Pressure-test scenario: Re-run the workflow on a pandas ETL job with one unsupported op and a dataset that does not fit GPU memory. 5. Success metric: The user gets a runnable GPU dataframe plan with explicit correctness and performance evidence paths.
<!-- PORTABILITY:START -->
Cross-Client Portability
This skill is written to stay usable across GitHub Copilot, Claude Code, Codex, and Gemini CLI.
- GitHub Copilot: keep the folder in a Copilot-visible skill or plugin path, or wrap the workflow as project instructions if the host does not support portable skill folders directly.
- Claude Code: keep the folder in a local skills directory or a compatible plugin or marketplace source.
- Codex: install or sync the folder into
$CODEX_HOME/skills/<skill-name>and restart Codex after major changes. - Gemini CLI: this repository generates a project command named
/skills:accelerated-computing-cudffrom this skill. Rebuild commands withpython scripts/export-gemini-skill.py accelerated-computing-cudfand then run/commands reloadinside Gemini CLI.
<!-- PORTABILITY:END -->
<!-- MCP:START -->
MCP Availability And Fallback
Preferred MCP Server: None required
- Fallback prompt: "Use the accelerated-computing-cudf skill without MCP. Rely on the local
SKILL.md, bundled references or scripts, and manual verification. Show the exact commands, evidence, and final checks you used before concluding." - If the current host does not expose a matching server, use the bundled references, scripts, native toolchain, and manual workflow already described in this skill.
- Treat direct local verification, rendered output, logs, tests, or screenshots as the fallback evidence path before completion.
<!-- MCP:END -->
Related Skills
- jupyter-notebook: Use it when the cuDF workflow also needs notebook-first experimentation or reproducible analysis.
- development-workflow: Use it when the migration also needs a scoped implementation and validation plan.
- documentation-verification: Use it when the deliverable is a guide or README that must keep commands and caveats accurate.
Evaluation Report
Evaluation of the accelerated-computing-cudf skill before publication through NVSkills-Eval.
This benchmark summarizes 3-Tier Evaluation from NVSkills-Eval results for the skill. The goal is to document whether the skill is safe, discoverable, effective, and useful for agents before it is published for broader workflow use.
Evaluation Summary
- Skill:
accelerated-computing-cudf - Evaluation date: 2026-05-29
- NVSkills-Eval profile:
external - Environment:
local - Dataset: 13 evaluation tasks
- Attempts per task: 2
- Pass threshold: 50%
- Overall verdict: PASS
Agents Used
claude-codecodex
Metrics Used
Reported benchmark dimensions:
- Security: checks whether skill-assisted execution avoids unsafe behavior such as secret leakage, destructive commands, or unauthorized access.
- Correctness: checks whether the agent follows the expected workflow and produces the correct final output.
- Discoverability: checks whether the agent loads the skill when relevant and avoids using it when irrelevant.
- Effectiveness: checks whether the agent performs measurably better with the skill than without it.
- Efficiency: checks whether the agent uses fewer tokens and avoids redundant work.
Underlying evaluation signals used in this run:
security(Security): checks for unsafe operations, secret leakage, and unauthorized access.skill_execution(Skill Execution): verifies that the agent loaded the expected skill and workflow.skill_efficiency(Efficiency): checks routing quality, decoy avoidance, and redundant tool usage.accuracy(Accuracy): grades final-answer correctness against the reference answer.goal_accuracy(Goal Accuracy): checks whether the overall user task completed successfully.behavior_check(Behavior Check): verifies expected behavior steps, including safety expectations.token_efficiency(Token Efficiency): compares token usage with and without the skill.
Test Tasks
The benchmark dataset contained 13 evaluation tasks:
- Positive tasks: 12 tasks where the skill was expected to activate.
- Negative tasks: 1 tasks where no skill was expected.
- Unlabeled tasks: 0 tasks where positive/negative intent could not be inferred.
Task composition is derived from the evaluation dataset when possible. Entries with expected_skill set are treated as positive skill-activation cases, while entries with expected_skill: null are treated as negative activation cases.
Results
| Dimension | Num | claude-code | codex |
|---|---|---|---|
| Security | 8 | 92% (+12%) | 100% (+0%) |
| Correctness | 8 | 96% (+10%) | 92% (+8%) |
| Discoverability | 8 | 84% (+26%) | 68% (+15%) |
| Effectiveness | 8 | 90% (+5%) | 86% (-0%) |
| Efficiency | 8 | 61% (+24%) | 50% (+10%) |
Score values show skill-assisted performance. Values in parentheses show uplift versus the no-skill baseline when baseline data is available.
Tier 1: Static Validation Summary
Tier 1 validation passed with observations. NVSkills-Eval ran 9 checks and found 7 total findings.
Top findings:
- MEDIUM SCHEMA/body_recommended_section: Missing recommended section: '## Instructions' (
skills/accelerated-computing-cudf/SKILL.md) - MEDIUM SCHEMA/body_recommended_section: Missing recommended section: '## Examples' (
skills/accelerated-computing-cudf/SKILL.md) - LOW QUALITY/quality_discoverability: Broad description without negative triggers may cause over-triggering (
skills/accelerated-computing-cudf/SKILL.md) - LOW QUALITY/quality_discoverability: No '## Purpose' section (
skills/accelerated-computing-cudf/SKILL.md) - LOW QUALITY/quality_reliability: No prerequisites/requirements documented (
skills/accelerated-computing-cudf/SKILL.md)
Tier 2: Deduplication Summary
Tier 2 validation passed. NVSkills-Eval ran 2 checks and found 0 total findings.
Notable observations:
- Context Deduplication: Collected 4 file(s)
- Inter-Skill Deduplication: Parsed skill 'accelerated-computing-cudf': 190 char description
Publication Recommendation
The skill is suitable to proceed toward NVSkills-Eval publication based on this benchmark. Skill owners should keep this file with the skill and refresh it when the evaluation dataset, skill behavior, or target agents materially change.
Changelog
All notable changes to the accelerated-computing-cudf skill will be documented in this file.
[2026-06-09] - Initial Import and Catalog Normalization
Added
- Imported
accelerated-computing-cudffromhttps://github.com/NVIDIA/skillsatskills/accelerated-computing-cudfpinned to129a1087a1853f32a950e2f7bbc0fd7d57b9d422. - Added repo-standard
Anti-Patterns,Verification Protocol, portability, MCP fallback, and related-skills sections. - Preserved the upstream benchmark, signature, skill card, and bundled references or scripts for provenance and later refreshes.
Changed
- Normalized
SKILL.mdfrontmatter to the shared catalog schema withversion: "1.2"andlast_updated: 2026-06-09. - Moved upstream-only top-level metadata into the nested
metadatablock so validation, export, and downstream sync stay consistent.
Fixed
- Aligned the imported skill with this repository's maintained-skill requirements and downstream sync workflow.
[
{
"id": "cudf-apply-udf__generic",
"question": "Task: Row-wise apply, applymap, and column-wise UDFs that should move to vectorized operations or Numba where appropriate\nTask folder: evals/files/cudf-apply-udf/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-apply-udf/code/generate_data.py",
"evals/files/cudf-apply-udf/code/udf_pipeline.py"
],
"ground_truth": "A successful answer uses the provided cudf-apply-udf starter files, especially code/udf_pipeline.py, to migrate the pandas DataFrame workload to cuDF where supported. It replaces row-wise apply/applymap or column UDF logic with vectorized cuDF expressions, Numba-compatible GPU logic, or a narrow compatibility boundary, preserves representative pandas results, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-csv-etl__generic",
"question": "Task: Basic CSV ETL pipeline \u2014 read, filter, compute columns, groupby aggregate, write to parquet\nTask folder: evals/files/cudf-csv-etl/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-csv-etl/code/etl_pipeline.py",
"evals/files/cudf-csv-etl/code/generate_data.py"
],
"ground_truth": "A successful answer uses the provided cudf-csv-etl starter files, especially code/etl_pipeline.py, to move CSV read, filtering, computed columns, groupby aggregation, and parquet output to cuDF. It preserves filter predicates, computed-column formulas, grouping keys, aggregate columns, generated data paths, output paths, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-groupby-agg__generic",
"question": "Task: Complex groupby with multiple agg functions, named aggregation, and transform\nTask folder: evals/files/cudf-groupby-agg/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-groupby-agg/code/generate_data.py",
"evals/files/cudf-groupby-agg/code/groupby_analysis.py"
],
"ground_truth": "A successful answer uses the provided cudf-groupby-agg starter files, especially code/groupby_analysis.py, to run the DataFrame loading and groupby work with cuDF. It preserves grouping keys, sum, mean, std, count, nunique, named aggregation, transform semantics or a documented compatibility boundary, output column names, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-multi-join__generic",
"question": "Task: Three-table join (orders, customers, products) with left/inner joins followed by aggregation\nTask folder: evals/files/cudf-multi-join/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-multi-join/code/generate_data.py",
"evals/files/cudf-multi-join/code/multi_join.py"
],
"ground_truth": "A successful answer uses the provided cudf-multi-join starter files, especially code/multi_join.py, to migrate the orders, customers, and products joins plus downstream filtering and aggregation to cuDF. It preserves left and inner join types, join keys, suffix behavior, row-count expectations, post-join filters, output schema, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-null-handling__generic",
"question": "Task: DataFrame with many nulls \u2014 fillna strategies, dropna, interpolate, isna masks, conditional fills\nTask folder: evals/files/cudf-null-handling/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-null-handling/code/generate_data.py",
"evals/files/cudf-null-handling/code/null_pipeline.py"
],
"ground_truth": "A successful answer uses the provided cudf-null-handling starter files, especially code/null_pipeline.py, to move null detection, fill, drop, mask, and conditional fill logic to cuDF where supported. It preserves scalar and dictionary fill rules, subset and threshold drop rules, NA-aware boolean masks, interpolation or other compatibility boundaries, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-parquet-io__generic",
"question": "Task: Read multiple parquet files, concatenate, filter, write partitioned output\nTask folder: evals/files/cudf-parquet-io/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-parquet-io/code/generate_data.py",
"evals/files/cudf-parquet-io/code/parquet_pipeline.py"
],
"ground_truth": "A successful answer uses the provided cudf-parquet-io starter files, especially code/parquet_pipeline.py, to migrate parquet reads, concatenation, filtering, column selection, dtype handling, and parquet writes to cuDF. It preserves multi-file input handling, partitioned output behavior, generated data paths, output paths, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-pivot-melt__generic",
"question": "Task: Pivot table creation, melt/unpivot, stack/unstack, and cross-tabulation\nTask folder: evals/files/cudf-pivot-melt/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-pivot-melt/code/generate_data.py",
"evals/files/cudf-pivot-melt/code/reshape_analysis.py"
],
"ground_truth": "A successful answer uses the provided cudf-pivot-melt starter files, especially code/reshape_analysis.py, to move supported reshape operations such as pivot, melt, stack/unstack, or crosstab-style logic to cuDF where practical. It preserves index labels, column labels, fill values, aggregation choices, output schema, compatibility boundaries, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-string-ops__generic",
"question": "Task: Text cleaning pipeline using pandas string accessor \u2014 lowercase, strip, regex extract, contains, replace\nTask folder: evals/files/cudf-string-ops/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-string-ops/code/clean_contacts.py",
"evals/files/cudf-string-ops/code/generate_data.py"
],
"ground_truth": "A successful answer uses the provided cudf-string-ops starter files, especially code/clean_contacts.py, to migrate string cleaning to cuDF string accessors for lowercase, strip, contains, replace, and extract-style operations. It preserves regex patterns, extracted columns, null handling, string dtype behavior, representative cleaned values, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-timeseries-resample__generic",
"question": "Task: Timestamped sensor data with resample to hourly/daily and rolling statistics\nTask folder: evals/files/cudf-timeseries-resample/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-timeseries-resample/code/generate_data.py",
"evals/files/cudf-timeseries-resample/code/timeseries_analysis.py"
],
"ground_truth": "A successful answer uses the provided cudf-timeseries-resample starter files, especially code/timeseries_analysis.py, to run datetime parsing, timestamp ordering, bucket creation, aggregation, and rolling computations with cuDF where supported. It preserves hourly and daily grouping semantics, missing buckets, rolling window sizes, output ordering, compatibility boundaries, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-window-functions__generic",
"question": "Task: Ranking, cumulative sums, rolling averages, expanding stats, and shift/lag operations\nTask folder: evals/files/cudf-window-functions/\nPrompt variant: generic\n\nUser prompt: Help me run this on the GPU\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-window-functions/code/generate_data.py",
"evals/files/cudf-window-functions/code/window_analysis.py"
],
"ground_truth": "A successful answer uses the provided cudf-window-functions starter files, especially code/window_analysis.py, to migrate ranking, cumulative operations, rolling calculations, expanding calculations, and shift/lag work to cuDF where supported. It preserves group keys, ordering columns, rank methods, window sizes, edge and null behavior, output names, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "source-cudf-null-fillna-semantics__generic",
"question": "Task: Preserve pandas nullable dtype and fillna semantics while migrating to cuDF.\nTask folder: evals/files/source-cudf-null-fillna-semantics/\nPrompt variant: generic\n\nUser prompt: Help me move this DataFrame cleanup to the GPU without messing up missing values.\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/source-cudf-null-fillna-semantics/NOTICE.md",
"evals/files/source-cudf-null-fillna-semantics/code/null_cleanup.py"
],
"ground_truth": "A successful answer uses the provided source-cudf-null-fillna-semantics starter files, especially code/null_cleanup.py, to migrate the cleanup workflow to cuDF without changing missing-value meaning. It preserves nullable integer, string, category-like, mask/where, fillna, and groupby semantics without lossy sentinel conversions, includes or describes pandas-versus-cuDF parity validation, and reports validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "cudf-native-stream-handoff-boundary__generic",
"question": "Task: Fix a threaded native GPU wrapper so cross-stream handoff and close/free ordering are correct.\nTask folder: evals/files/cudf-native-stream-handoff-boundary/\nPrompt variant: generic\n\nUser prompt: This threaded GPU wrapper sometimes returns stale checksums after one\nworker hands a device buffer to another. Can you make the handoff correct\nwithout blocking the whole device on every transfer, and keep cleanup safe\nfor queued GPU work?\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": "accelerated-computing-cudf",
"expected_script": null,
"files": [
"evals/files/cudf-native-stream-handoff-boundary/NOTICE.md",
"evals/files/cudf-native-stream-handoff-boundary/code/run_smoke.sh",
"evals/files/cudf-native-stream-handoff-boundary/code/threaded_handoff.cu"
],
"ground_truth": "A successful answer uses the provided cudf-native-stream-handoff-boundary starter files, especially code/threaded_handoff.cu, to fix cross-thread or cross-stream GPU handoff by tying CUDA event readiness to the object dependency. It orders consumer work after producer writes, orders destruction or free after last stream use, preserves asynchronous overlap where practical, and reports compile or smoke validation performed or the runtime blocker.",
"expected_behavior": []
},
{
"id": "negative-deep-learning-training__generic",
"question": "Task: Assess whether a PyTorch training performance issue belongs in NVIDIA GPU data science migration guidance.\nTask folder: evals/files/negative-deep-learning-training/\nPrompt variant: generic\n\nUser prompt: This PyTorch training script underutilizes my H100. Help me speed up model\ntraining on the GPU.\n\nUse the provided starter workspace for this task. Modify the starter file(s) under the provided `code/` directory. Run the relevant smoke or validation command from that workspace when practical, and report the changed files and validation result.",
"expected_skill": null,
"expected_script": null,
"files": [
"evals/files/negative-deep-learning-training/code/train.py"
],
"ground_truth": "A successful answer treats the provided train.py context as a PyTorch/deep-learning training performance task rather than a cuDF migration. It keeps guidance focused on model training, data loading, batching, mixed precision, profiling, or other training-specific tactics, and only mentions cuDF as optional upstream tabular ETL when that is directly relevant.",
"expected_behavior": []
}
]
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate synthetic insurance claims data for UDF processing."""
import os
import numpy as np
import pandas as pd
SEED = 42
N_ROWS = 40_000
def generate():
if os.path.exists("claims.csv"):
return
rng = np.random.default_rng(SEED)
policy_types = ["auto", "home", "health", "life", "travel"]
risk_levels = ["low", "medium", "high"]
regions = ["northeast", "southeast", "midwest", "west", "pacific"]
df = pd.DataFrame({
"claim_id": range(N_ROWS),
"policy_type": rng.choice(policy_types, N_ROWS),
"risk_level": rng.choice(risk_levels, N_ROWS, p=[0.5, 0.35, 0.15]),
"region": rng.choice(regions, N_ROWS),
"age": rng.integers(18, 85, N_ROWS),
"claim_amount": np.round(rng.exponential(5000, N_ROWS), 2),
"deductible": np.round(rng.choice([250, 500, 1000, 2000, 5000], N_ROWS).astype(float), 2),
"premium_monthly": np.round(rng.uniform(50, 800, N_ROWS), 2),
"years_as_customer": rng.integers(0, 30, N_ROWS),
"num_prior_claims": rng.integers(0, 10, N_ROWS),
"credit_score": rng.integers(300, 850, N_ROWS),
"property_value": np.round(rng.uniform(50_000, 1_000_000, N_ROWS), 2),
})
df.to_csv("claims.csv", index=False)
print(f"Generated {len(df)} insurance claims -> claims.csv")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""UDF-heavy processing pipeline on insurance claims data.
Uses apply(), applymap(), and custom functions for row-wise and
element-wise transformations on a pandas DataFrame.
"""
import numpy as np
import pandas as pd
from generate_data import generate
def load_data():
generate()
df = pd.read_csv("claims.csv")
print(f"Loaded {len(df)} claims")
return df
# --- Row-wise UDFs used with apply(axis=1) ---
def calculate_risk_score(row):
"""Complex row-wise risk scoring function."""
base_score = 50
# Age factor
if row["age"] < 25:
base_score += 15
elif row["age"] > 65:
base_score += 10
else:
base_score -= 5
# Claims history
base_score += row["num_prior_claims"] * 8
# Credit score factor
if row["credit_score"] >= 750:
base_score -= 20
elif row["credit_score"] >= 650:
base_score -= 10
elif row["credit_score"] < 550:
base_score += 15
# Risk level multiplier
if row["risk_level"] == "high":
base_score *= 1.5
elif row["risk_level"] == "medium":
base_score *= 1.2
# Loyalty discount
if row["years_as_customer"] > 10:
base_score *= 0.85
elif row["years_as_customer"] > 5:
base_score *= 0.92
return round(base_score, 2)
def calculate_payout(row):
"""Calculate adjusted payout amount based on multiple conditions."""
amount = row["claim_amount"]
deductible = row["deductible"]
net = max(0, amount - deductible)
# Cap by policy type
caps = {"auto": 50_000, "home": 200_000, "health": 100_000,
"life": 500_000, "travel": 10_000}
cap = caps.get(row["policy_type"], 50_000)
net = min(net, cap)
# Loyalty bonus: extra 5% for long-term customers
if row["years_as_customer"] > 15:
net *= 1.05
# High-risk penalty: reduce by 10%
if row["risk_level"] == "high" and row["num_prior_claims"] > 5:
net *= 0.90
return round(net, 2)
def classify_claim_tier(row):
"""Classify claim into processing tier based on multiple factors."""
amount = row["claim_amount"]
risk = row["risk_level"]
priors = row["num_prior_claims"]
if amount > 20_000 or (risk == "high" and priors > 3):
return "tier_3_manual"
elif amount > 5_000 or (risk == "medium" and priors > 2):
return "tier_2_review"
else:
return "tier_1_auto"
# --- Column-wise UDFs ---
def normalize_score(series):
"""Min-max normalize a numeric series."""
return (series - series.min()) / (series.max() - series.min())
def winsorize(series, lower=0.05, upper=0.95):
"""Clip values at the given percentiles."""
lo = series.quantile(lower)
hi = series.quantile(upper)
return series.clip(lo, hi)
# --- Element-wise UDF ---
def format_currency(val):
"""Format a numeric value as currency string."""
if pd.isna(val):
return "$0.00"
return f"${val:,.2f}"
def credit_bucket(val):
"""Bucket a credit score into a category."""
if val >= 750:
return "excellent"
elif val >= 700:
return "good"
elif val >= 650:
return "fair"
elif val >= 550:
return "poor"
else:
return "very_poor"
def process_claims(df):
"""Apply all UDFs to the claims DataFrame."""
# Row-wise apply (the expensive operations)
print("Computing risk scores (row-wise apply)...")
df["risk_score"] = df.apply(calculate_risk_score, axis=1)
print("Computing payouts (row-wise apply)...")
df["payout"] = df.apply(calculate_payout, axis=1)
print("Classifying claims (row-wise apply)...")
df["claim_tier"] = df.apply(classify_claim_tier, axis=1)
# Column-wise UDFs
print("Normalizing and winsorizing...")
df["risk_score_norm"] = normalize_score(df["risk_score"])
df["claim_amount_winsorized"] = winsorize(df["claim_amount"])
df["premium_norm"] = normalize_score(df["premium_monthly"])
# Element-wise apply (applymap-style via apply on columns)
print("Formatting and bucketing...")
df["credit_bucket"] = df["credit_score"].apply(credit_bucket)
df["payout_formatted"] = df["payout"].apply(format_currency)
# Element-wise on multiple numeric columns
numeric_cols = ["claim_amount", "deductible", "premium_monthly", "property_value"]
formatted = df[numeric_cols].applymap(format_currency)
for col in numeric_cols:
df[f"{col}_fmt"] = formatted[col]
return df
def summarize(df):
"""Summarize processed claims."""
print(f"\nProcessed {len(df)} claims")
print(f"Risk score stats: mean={df['risk_score'].mean():.1f}, "
f"std={df['risk_score'].std():.1f}")
print(f"Total payouts: ${df['payout'].sum():,.2f}")
tier_counts = df["claim_tier"].value_counts()
print(f"\nClaim tiers:\n{tier_counts}")
credit_dist = df["credit_bucket"].value_counts()
print(f"\nCredit distribution:\n{credit_dist}")
by_type = df.groupby("policy_type").agg(
avg_risk=("risk_score", "mean"),
total_payout=("payout", "sum"),
claim_count=("claim_id", "count"),
).round(2)
print(f"\nBy policy type:\n{by_type}")
def main():
df = load_data()
df = process_claims(df)
summarize(df)
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""CSV ETL pipeline: read, filter, compute, groupby, write parquet.
Reads sales.csv, filters to completed orders, adds computed columns
(revenue, discounted_revenue, age_group), runs a groupby aggregation
by region and product, and writes the summary to parquet.
"""
import numpy as np
import pandas as pd
from generate_data import generate
def load_data():
generate()
df = pd.read_csv("sales.csv")
print(f"Loaded {len(df)} rows from sales.csv")
return df
def filter_completed(df):
"""Keep only completed orders with quantity >= 2."""
mask = (df["status"] == "completed") & (df["quantity"] >= 2)
filtered = df[mask].copy()
print(f"Filtered to {len(filtered)} completed orders")
return filtered
def add_computed_columns(df):
"""Add revenue, discounted revenue, and age group columns."""
df["revenue"] = df["quantity"] * df["unit_price"]
df["discounted_revenue"] = df["revenue"] * (1 - df["discount_pct"])
bins = [0, 25, 35, 50, 65, 100]
labels = ["18-25", "26-35", "36-50", "51-65", "65+"]
df["age_group"] = pd.cut(df["customer_age"], bins=bins, labels=labels)
df["high_value"] = (df["discounted_revenue"] > 500).astype(int)
print(f"Added computed columns; {df['high_value'].sum()} high-value orders")
return df
def aggregate_by_region_product(df):
"""Groupby region + product, compute summary statistics."""
summary = (
df.groupby(["region", "product"])
.agg(
total_revenue=("revenue", "sum"),
total_discounted=("discounted_revenue", "sum"),
order_count=("order_id", "count"),
avg_quantity=("quantity", "mean"),
avg_unit_price=("unit_price", "mean"),
high_value_count=("high_value", "sum"),
)
.reset_index()
)
summary["avg_discount_impact"] = (
1 - summary["total_discounted"] / summary["total_revenue"]
)
summary = summary.sort_values("total_revenue", ascending=False)
print(f"Aggregated into {len(summary)} region-product groups")
return summary
def write_output(summary):
"""Write the summary to a parquet file."""
summary.to_parquet("sales_summary.parquet", index=False)
print("Wrote sales_summary.parquet")
def main():
df = load_data()
df = filter_completed(df)
df = add_computed_columns(df)
summary = aggregate_by_region_product(df)
write_output(summary)
print("\nTop 5 region-product combos by revenue:")
print(summary.head(5).to_string(index=False))
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate a synthetic sales CSV for the ETL pipeline."""
import os
import numpy as np
import pandas as pd
SEED = 42
N_ROWS = 50_000
def generate():
if os.path.exists("sales.csv"):
return
rng = np.random.default_rng(SEED)
regions = ["North", "South", "East", "West"]
products = ["Widget", "Gadget", "Doohickey", "Thingamajig", "Whatchamacallit"]
statuses = ["completed", "pending", "returned", "cancelled"]
df = pd.DataFrame({
"order_id": range(N_ROWS),
"region": rng.choice(regions, N_ROWS),
"product": rng.choice(products, N_ROWS),
"quantity": rng.integers(1, 50, N_ROWS),
"unit_price": np.round(rng.uniform(5.0, 500.0, N_ROWS), 2),
"discount_pct": np.round(rng.uniform(0.0, 0.3, N_ROWS), 3),
"status": rng.choice(statuses, N_ROWS, p=[0.7, 0.1, 0.1, 0.1]),
"customer_age": rng.integers(18, 80, N_ROWS),
})
df.to_csv("sales.csv", index=False)
print(f"Generated {len(df)} rows -> sales.csv")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate synthetic employee performance data."""
import os
import numpy as np
import pandas as pd
SEED = 42
N_EMPLOYEES = 50_000
def generate():
if os.path.exists("employees.csv"):
return
rng = np.random.default_rng(SEED)
departments = ["Engineering", "Sales", "Marketing", "Finance", "HR", "Operations"]
levels = ["Junior", "Mid", "Senior", "Lead", "Principal"]
offices = ["NYC", "SF", "London", "Berlin", "Tokyo", "Sydney"]
df = pd.DataFrame({
"employee_id": range(N_EMPLOYEES),
"department": rng.choice(departments, N_EMPLOYEES),
"level": rng.choice(levels, N_EMPLOYEES, p=[0.3, 0.3, 0.2, 0.12, 0.08]),
"office": rng.choice(offices, N_EMPLOYEES),
"salary": np.round(rng.normal(85_000, 25_000, N_EMPLOYEES).clip(30_000, 300_000), 2),
"bonus": np.round(rng.exponential(5_000, N_EMPLOYEES), 2),
"performance_score": np.round(rng.normal(3.5, 0.8, N_EMPLOYEES).clip(1.0, 5.0), 2),
"years_tenure": rng.integers(0, 25, N_EMPLOYEES),
"projects_completed": rng.integers(0, 50, N_EMPLOYEES),
"training_hours": np.round(rng.exponential(20, N_EMPLOYEES), 1),
})
df.to_csv("employees.csv", index=False)
print(f"Generated {len(df)} employee records -> employees.csv")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Complex groupby aggregation and transform pipeline.
Performs department-level, multi-key groupby, named aggregation,
and transform-based feature engineering on employee data.
"""
import numpy as np
import pandas as pd
from generate_data import generate
def load_data():
generate()
df = pd.read_csv("employees.csv")
print(f"Loaded {len(df)} employees")
return df
def department_summary(df):
"""Basic department-level aggregation with multiple functions."""
dept = df.groupby("department").agg(
headcount=("employee_id", "count"),
avg_salary=("salary", "mean"),
median_salary=("salary", "median"),
std_salary=("salary", "std"),
total_bonus=("bonus", "sum"),
avg_perf=("performance_score", "mean"),
unique_levels=("level", "nunique"),
unique_offices=("office", "nunique"),
avg_tenure=("years_tenure", "mean"),
total_projects=("projects_completed", "sum"),
).reset_index()
dept = dept.sort_values("avg_salary", ascending=False)
print(f"Department summary: {len(dept)} departments")
return dept
def multi_key_aggregation(df):
"""Groupby on department + level with named aggregation."""
result = df.groupby(["department", "level"]).agg(
count=("employee_id", "count"),
salary_mean=("salary", "mean"),
salary_min=("salary", "min"),
salary_max=("salary", "max"),
salary_sum=("salary", "sum"),
bonus_mean=("bonus", "mean"),
perf_mean=("performance_score", "mean"),
perf_std=("performance_score", "std"),
tenure_mean=("years_tenure", "mean"),
projects_sum=("projects_completed", "sum"),
).reset_index()
result["salary_range"] = result["salary_max"] - result["salary_min"]
print(f"Multi-key aggregation: {len(result)} groups")
return result
def office_department_crosstab(df):
"""Three-key groupby: department + office + level."""
cross = df.groupby(["department", "office", "level"]).agg(
headcount=("employee_id", "count"),
avg_salary=("salary", "mean"),
total_training=("training_hours", "sum"),
).reset_index()
print(f"Cross-tab: {len(cross)} groups")
return cross
def add_transform_features(df):
"""Use groupby transform to add group-relative features."""
# Department-level transforms
df["dept_avg_salary"] = df.groupby("department")["salary"].transform("mean")
df["dept_std_salary"] = df.groupby("department")["salary"].transform("std")
df["salary_zscore"] = (df["salary"] - df["dept_avg_salary"]) / df["dept_std_salary"]
# Level-level transforms
df["level_avg_perf"] = df.groupby("level")["performance_score"].transform("mean")
df["perf_vs_level"] = df["performance_score"] - df["level_avg_perf"]
# Department rank by salary
df["dept_salary_rank"] = df.groupby("department")["salary"].rank(
method="dense", ascending=False
)
# Department + level cumulative count
df["dept_level_count"] = df.groupby(["department", "level"]).cumcount() + 1
# Percent of department total
df["dept_salary_total"] = df.groupby("department")["salary"].transform("sum")
df["salary_pct_of_dept"] = df["salary"] / df["dept_salary_total"]
outlier_count = (df["salary_zscore"].abs() > 2).sum()
print(f"Transform features added; {outlier_count} salary outliers (|z| > 2)")
return df
def top_performers_per_dept(df):
"""Get top 5 performers per department using groupby + nlargest."""
top = (
df.groupby("department")
.apply(lambda g: g.nlargest(5, "performance_score"))
.reset_index(drop=True)
)
print(f"Top performers: {len(top)} rows")
return top
def main():
df = load_data()
dept_summary = department_summary(df)
multi_key = multi_key_aggregation(df)
cross = office_department_crosstab(df)
df_with_transforms = add_transform_features(df)
top_perf = top_performers_per_dept(df)
print(f"\nDepartment summary:\n{dept_summary.to_string(index=False)}")
print(f"\nSample transformed rows:\n"
f"{df_with_transforms[['department', 'level', 'salary', 'salary_zscore', 'perf_vs_level', 'dept_salary_rank']].head(10).to_string(index=False)}")
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate three related CSVs: orders, customers, products."""
import os
import numpy as np
import pandas as pd
SEED = 42
N_CUSTOMERS = 3_000
N_PRODUCTS = 200
N_ORDERS = 80_000
def generate():
if os.path.exists("orders.csv"):
return
rng = np.random.default_rng(SEED)
# --- customers ---
tiers = ["bronze", "silver", "gold", "platinum"]
customers = pd.DataFrame({
"customer_id": range(N_CUSTOMERS),
"customer_name": [f"Cust_{i:05d}" for i in range(N_CUSTOMERS)],
"tier": rng.choice(tiers, N_CUSTOMERS, p=[0.4, 0.3, 0.2, 0.1]),
"country": rng.choice(["US", "UK", "DE", "JP", "BR", "IN"], N_CUSTOMERS),
"credit_limit": np.round(rng.uniform(500, 50_000, N_CUSTOMERS), 2),
})
# --- products ---
categories = ["electronics", "clothing", "food", "tools", "toys"]
products = pd.DataFrame({
"product_id": range(N_PRODUCTS),
"product_name": [f"Prod_{i:04d}" for i in range(N_PRODUCTS)],
"category": rng.choice(categories, N_PRODUCTS),
"base_price": np.round(rng.uniform(2.0, 800.0, N_PRODUCTS), 2),
"weight_kg": np.round(rng.uniform(0.1, 30.0, N_PRODUCTS), 2),
})
# --- orders (some customer_ids intentionally out of range to test left join) ---
orders = pd.DataFrame({
"order_id": range(N_ORDERS),
"customer_id": rng.integers(0, N_CUSTOMERS + 200, N_ORDERS),
"product_id": rng.integers(0, N_PRODUCTS, N_ORDERS),
"quantity": rng.integers(1, 20, N_ORDERS),
"order_total": np.round(rng.uniform(5.0, 2000.0, N_ORDERS), 2),
"channel": rng.choice(["web", "mobile", "store", "phone"], N_ORDERS),
})
customers.to_csv("customers.csv", index=False)
products.to_csv("products.csv", index=False)
orders.to_csv("orders.csv", index=False)
print(f"Generated {N_CUSTOMERS} customers, {N_PRODUCTS} products, {N_ORDERS} orders")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Three-table join pipeline with aggregation.
Joins orders with customers (left join) and products (inner join),
then computes per-customer and per-category summaries.
"""
import numpy as np
import pandas as pd
from generate_data import generate
def load_tables():
generate()
orders = pd.read_csv("orders.csv")
customers = pd.read_csv("customers.csv")
products = pd.read_csv("products.csv")
print(f"Loaded orders={len(orders)}, customers={len(customers)}, products={len(products)}")
return orders, customers, products
def join_tables(orders, customers, products):
"""Left-join orders->customers, then inner-join with products."""
# Left join: keep all orders even if customer_id is missing
merged = orders.merge(customers, on="customer_id", how="left")
print(f"After left join with customers: {len(merged)} rows, "
f"{merged['customer_name'].isna().sum()} unmatched customers")
# Inner join: drop orders whose product_id doesn't match
merged = merged.merge(products, on="product_id", how="inner")
print(f"After inner join with products: {len(merged)} rows")
# Computed columns
merged["line_total"] = merged["quantity"] * merged["base_price"]
merged["total_weight"] = merged["quantity"] * merged["weight_kg"]
merged["over_credit"] = (merged["order_total"] > merged["credit_limit"]).fillna(False)
return merged
def customer_summary(merged):
"""Per-customer aggregation."""
cust = (
merged.groupby("customer_id")
.agg(
num_orders=("order_id", "count"),
total_spent=("order_total", "sum"),
avg_order=("order_total", "mean"),
unique_products=("product_id", "nunique"),
total_weight=("total_weight", "sum"),
times_over_credit=("over_credit", "sum"),
tier=("tier", "first"),
country=("country", "first"),
)
.reset_index()
.sort_values("total_spent", ascending=False)
)
print(f"Customer summary: {len(cust)} customers")
return cust
def category_summary(merged):
"""Per-category aggregation."""
cat = (
merged.groupby("category")
.agg(
num_orders=("order_id", "count"),
total_revenue=("line_total", "sum"),
avg_quantity=("quantity", "mean"),
unique_customers=("customer_id", "nunique"),
avg_weight=("total_weight", "mean"),
)
.reset_index()
.sort_values("total_revenue", ascending=False)
)
print(f"Category summary: {len(cat)} categories")
return cat
def tier_channel_summary(merged):
"""Cross-tabulation of tier x channel."""
cross = (
merged.groupby(["tier", "channel"])
.agg(
order_count=("order_id", "count"),
revenue=("line_total", "sum"),
)
.reset_index()
)
# Pivot to wide format
pivot = cross.pivot_table(
index="tier", columns="channel", values="revenue",
aggfunc="sum", fill_value=0,
)
print(f"Tier-channel pivot:\n{pivot}")
return cross
def main():
orders, customers, products = load_tables()
merged = join_tables(orders, customers, products)
cust_summary = customer_summary(merged)
cat_summary = category_summary(merged)
tier_ch = tier_channel_summary(merged)
print(f"\nTop 5 customers by spend:\n{cust_summary.head(5).to_string(index=False)}")
print(f"\nCategory breakdown:\n{cat_summary.to_string(index=False)}")
if __name__ == "__main__":
main()
#!/usr/bin/env bash
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
set -euo pipefail
script_dir="$(CDPATH= cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd)"
tmp_dir="$(mktemp -d "${TMPDIR:-/var/tmp}/threaded-handoff.XXXXXX")"
trap 'rm -rf "$tmp_dir"' EXIT
nvcc -std=c++17 -O2 "$script_dir/threaded_handoff.cu" -o "$tmp_dir/threaded_handoff"
"$tmp_dir/threaded_handoff"
/*
* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
* SPDX-License-Identifier: Apache-2.0
*/
#include <cuda_runtime.h>
#include <cstdint>
#include <cstdio>
#include <cstdlib>
#include <memory>
namespace {
constexpr int kRows = 1 << 20;
constexpr int kValue = 7;
constexpr int kTrials = 32;
void check_cuda(cudaError_t status, const char* what)
{
if (status != cudaSuccess) {
std::fprintf(stderr, "%s failed: %s\n", what, cudaGetErrorString(status));
std::abort();
}
}
__global__ void fill_kernel(int* data, int rows, int value)
{
int idx = blockIdx.x * blockDim.x + threadIdx.x;
if (idx >= rows) return;
int adjusted = value;
for (int spin = 0; spin < 4096; ++spin) {
adjusted += (spin & 1);
adjusted -= (spin & 1);
}
data[idx] = adjusted;
}
__global__ void checksum_kernel(const int* data, std::uint64_t* out, int rows)
{
int idx = blockIdx.x * blockDim.x + threadIdx.x;
if (idx < rows) {
atomicAdd(reinterpret_cast<unsigned long long*>(out),
static_cast<unsigned long long>(data[idx]));
}
}
struct NativeGpuTable {
int* data{};
int rows{};
cudaStream_t producer_stream{};
explicit NativeGpuTable(int row_count) : rows(row_count)
{
check_cuda(cudaStreamCreateWithFlags(&producer_stream, cudaStreamNonBlocking),
"create producer stream");
check_cuda(cudaMalloc(&data, sizeof(int) * rows), "allocate table data");
}
~NativeGpuTable()
{
if (data != nullptr) {
cudaFree(data);
}
if (producer_stream != nullptr) {
cudaStreamDestroy(producer_stream);
}
}
};
std::shared_ptr<NativeGpuTable> build_table_async(int value)
{
auto table = std::make_shared<NativeGpuTable>(kRows);
int block = 256;
int grid = (table->rows + block - 1) / block;
fill_kernel<<<grid, block, 0, table->producer_stream>>>(table->data, table->rows, value);
check_cuda(cudaGetLastError(), "launch fill kernel");
return table;
}
std::uint64_t consume_on_stream(const std::shared_ptr<NativeGpuTable>& table,
cudaStream_t consumer_stream)
{
std::uint64_t* d_sum{};
std::uint64_t h_sum{};
int block = 256;
int grid = (table->rows + block - 1) / block;
check_cuda(cudaMalloc(&d_sum, sizeof(std::uint64_t)), "allocate checksum");
check_cuda(cudaMemsetAsync(d_sum, 0, sizeof(std::uint64_t), consumer_stream),
"clear checksum");
checksum_kernel<<<grid, block, 0, consumer_stream>>>(table->data, d_sum, table->rows);
check_cuda(cudaGetLastError(), "launch checksum kernel");
check_cuda(cudaMemcpyAsync(&h_sum,
d_sum,
sizeof(std::uint64_t),
cudaMemcpyDeviceToHost,
consumer_stream),
"copy checksum");
check_cuda(cudaStreamSynchronize(consumer_stream), "sync consumer stream");
check_cuda(cudaFree(d_sum), "free checksum");
return h_sum;
}
} // namespace
int main()
{
cudaStream_t consumer_stream{};
check_cuda(cudaStreamCreateWithFlags(&consumer_stream, cudaStreamNonBlocking),
"create consumer stream");
std::uint64_t expected = static_cast<std::uint64_t>(kRows) * kValue;
int failures = 0;
for (int trial = 0; trial < kTrials; ++trial) {
auto table = build_table_async(kValue);
std::uint64_t actual = consume_on_stream(table, consumer_stream);
if (actual != expected) {
std::fprintf(stderr,
"trial %d checksum mismatch: got %llu expected %llu\n",
trial,
static_cast<unsigned long long>(actual),
static_cast<unsigned long long>(expected));
++failures;
}
}
check_cuda(cudaStreamDestroy(consumer_stream), "destroy consumer stream");
if (failures != 0) {
std::fprintf(stderr, "%d stale handoff checks observed\n", failures);
return 1;
}
std::puts("all handoffs matched expected checksum");
return 0;
}
Notice
This task is an original synthetic fixture. No upstream source code was copied.
It is inspired by public CUDA stream/event ordering guidance and public cuDF native/JVM wrapper concepts. The starter program is intentionally small so the task focuses on object readiness, cross-stream consumption, and device-memory lifetime ordering.
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate synthetic data with intentional null patterns."""
import os
import numpy as np
import pandas as pd
SEED = 42
N_ROWS = 40_000
def generate():
if os.path.exists("messy_data.csv"):
return
rng = np.random.default_rng(SEED)
df = pd.DataFrame({
"id": range(N_ROWS),
"group": rng.choice(["A", "B", "C", "D"], N_ROWS),
"temperature": rng.normal(22.0, 3.0, N_ROWS),
"humidity": rng.uniform(20, 90, N_ROWS),
"pressure": rng.normal(1013, 5, N_ROWS),
"wind_speed": rng.exponential(10, N_ROWS),
"visibility": rng.uniform(1, 30, N_ROWS),
"uv_index": rng.integers(0, 12, N_ROWS).astype(float),
"air_quality": rng.choice(["good", "moderate", "poor", "hazardous"], N_ROWS),
"station_code": rng.choice(["ST01", "ST02", "ST03", "ST04", "ST05"], N_ROWS),
})
# Introduce nulls with different patterns
# Random scattered nulls (~15% each)
for col in ["temperature", "humidity", "pressure"]:
mask = rng.random(N_ROWS) < 0.15
df.loc[mask, col] = np.nan
# Block nulls (sensor offline for stretches)
for start in [5000, 15000, 28000]:
df.loc[start:start + 500, "wind_speed"] = np.nan
df.loc[start:start + 300, "visibility"] = np.nan
# Correlated nulls (uv_index missing when visibility is low)
low_vis = df["visibility"] < 5
df.loc[low_vis & (rng.random(N_ROWS) < 0.7), "uv_index"] = np.nan
# String column nulls
str_mask = rng.random(N_ROWS) < 0.10
df.loc[str_mask, "air_quality"] = np.nan
df["temperature"] = df["temperature"].round(2)
df["humidity"] = df["humidity"].round(1)
df["pressure"] = df["pressure"].round(1)
df["wind_speed"] = df["wind_speed"].round(2)
df["visibility"] = df["visibility"].round(1)
df.to_csv("messy_data.csv", index=False)
null_pcts = df.isnull().mean() * 100
print(f"Generated {len(df)} rows with null percentages:\n{null_pcts.to_string()}")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Null handling pipeline: detect, fill, drop, interpolate, and report.
Demonstrates various pandas null-handling strategies on messy weather data.
"""
import numpy as np
import pandas as pd
from generate_data import generate
def load_data():
generate()
df = pd.read_csv("messy_data.csv")
print(f"Loaded {len(df)} rows")
print(f"Null counts:\n{df.isnull().sum()}")
return df
def analyze_nulls(df):
"""Build a null analysis report."""
null_counts = df.isnull().sum()
null_pcts = df.isnull().mean() * 100
report = pd.DataFrame({
"null_count": null_counts,
"null_pct": null_pcts.round(2),
"dtype": df.dtypes,
})
# Per-group null rates
group_nulls = df.groupby("group").apply(
lambda g: g.isnull().sum()
).reset_index()
print(f"Null report:\n{report}")
return report, group_nulls
def fill_with_strategies(df):
"""Apply different fill strategies to different columns."""
filled = df.copy()
# Scalar fill
filled["uv_index"] = filled["uv_index"].fillna(0)
# Dict fill (different values per column)
filled = filled.fillna({
"air_quality": "unknown",
"visibility": filled["visibility"].median(),
})
# Forward fill for block-missing wind data
filled["wind_speed"] = filled["wind_speed"].ffill()
# Backward fill for any remaining at the start
filled["wind_speed"] = filled["wind_speed"].bfill()
# Group-specific mean fill for temperature
group_means = df.groupby("group")["temperature"].transform("mean")
filled["temperature"] = filled["temperature"].fillna(group_means)
# Conditional fill: humidity depends on air_quality
quality_median = df.groupby("air_quality")["humidity"].median()
for quality, median_val in quality_median.items():
mask = filled["humidity"].isna() & (filled["air_quality"] == quality)
filled.loc[mask, "humidity"] = median_val
# Fill remaining humidity nulls with global median
filled["humidity"] = filled["humidity"].fillna(filled["humidity"].median())
print(f"After fills, remaining nulls:\n{filled.isnull().sum()}")
return filled
def interpolate_pressure(df):
"""Interpolate pressure readings within each station."""
interp_frames = []
for station, group in df.groupby("station_code"):
g = group.copy()
g["pressure"] = g["pressure"].interpolate(method="linear", limit=10)
g["pressure"] = g["pressure"].bfill().ffill()
interp_frames.append(g)
result = pd.concat(interp_frames, ignore_index=True)
remaining = result["pressure"].isna().sum()
print(f"After interpolation, {remaining} pressure nulls remain")
return result
def dropna_analysis(df):
"""Demonstrate dropna with various parameters."""
# Drop rows where all numeric columns are null
numeric_cols = df.select_dtypes(include=[np.number]).columns.tolist()
dropped_all = df.dropna(subset=numeric_cols, how="all")
print(f"dropna(how='all') on numeric: {len(df)} -> {len(dropped_all)}")
# Drop rows where more than 3 columns are null
dropped_thresh = df.dropna(thresh=len(df.columns) - 3)
print(f"dropna(thresh={len(df.columns) - 3}): {len(df)} -> {len(dropped_thresh)}")
# Drop rows with any null in key columns
key_cols = ["temperature", "humidity", "pressure"]
dropped_subset = df.dropna(subset=key_cols)
print(f"dropna(subset={key_cols}): {len(df)} -> {len(dropped_subset)}")
return dropped_subset
def create_null_indicators(df):
"""Create boolean indicator columns for null patterns."""
indicator_cols = ["temperature", "humidity", "pressure", "wind_speed", "uv_index"]
for col in indicator_cols:
df[f"{col}_missing"] = df[col].isna().astype(int)
df["total_missing"] = df[[f"{c}_missing" for c in indicator_cols]].sum(axis=1)
df["has_any_missing"] = (df["total_missing"] > 0).astype(int)
# Null pattern string
df["null_pattern"] = ""
for col in indicator_cols:
df["null_pattern"] = df["null_pattern"] + df[f"{col}_missing"].astype(str)
pattern_counts = df["null_pattern"].value_counts().head(10)
print(f"\nTop null patterns:\n{pattern_counts}")
return df
def main():
df = load_data()
report, group_nulls = analyze_nulls(df)
df_with_indicators = create_null_indicators(df)
dropped = dropna_analysis(df)
filled = fill_with_strategies(df)
result = interpolate_pressure(filled)
print(f"\nFinal null check:\n{result.isnull().sum()}")
print(f"\nSample rows:\n{result.head(5).to_string(index=False)}")
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate multiple parquet files simulating partitioned log data."""
import os
import numpy as np
import pandas as pd
from pathlib import Path
SEED = 42
N_PER_FILE = 10_000
N_FILES = 6
def generate():
outdir = Path("raw_logs")
if outdir.exists() and len(list(outdir.glob("*.parquet"))) == N_FILES:
return
outdir.mkdir(exist_ok=True)
rng = np.random.default_rng(SEED)
endpoints = ["/api/users", "/api/orders", "/api/products",
"/api/health", "/api/search", "/api/auth"]
methods = ["GET", "POST", "PUT", "DELETE"]
status_codes = [200, 201, 204, 301, 400, 401, 403, 404, 500, 502, 503]
status_weights = [0.50, 0.10, 0.05, 0.02, 0.08, 0.05, 0.03, 0.07, 0.04, 0.03, 0.03]
regions = ["us-east-1", "us-west-2", "eu-west-1", "ap-southeast-1"]
for i in range(N_FILES):
base_date = pd.Timestamp("2024-01-01") + pd.Timedelta(days=i * 5)
timestamps = base_date + pd.to_timedelta(
rng.integers(0, 5 * 86400, N_PER_FILE), unit="s"
)
df = pd.DataFrame({
"timestamp": timestamps,
"endpoint": rng.choice(endpoints, N_PER_FILE),
"method": rng.choice(methods, N_PER_FILE, p=[0.6, 0.2, 0.1, 0.1]),
"status_code": rng.choice(status_codes, N_PER_FILE, p=status_weights),
"response_time_ms": np.round(rng.exponential(150, N_PER_FILE), 2),
"bytes_sent": rng.integers(100, 50_000, N_PER_FILE),
"user_id": rng.integers(1, 5_000, N_PER_FILE),
"region": rng.choice(regions, N_PER_FILE),
"is_cached": rng.choice([True, False], N_PER_FILE, p=[0.3, 0.7]),
})
fname = outdir / f"logs_batch_{i:03d}.parquet"
df.to_parquet(fname, index=False)
print(f"Wrote {fname} ({len(df)} rows)")
print(f"Generated {N_FILES} parquet files in {outdir}/")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Parquet I/O pipeline: read multiple files, concatenate, filter, write partitioned.
Reads log data from multiple parquet files, concatenates them, applies
filters and transformations, then writes partitioned parquet output.
"""
import os
import numpy as np
import pandas as pd
from pathlib import Path
from generate_data import generate
def load_all_parquet(input_dir):
"""Read all parquet files from a directory and concatenate."""
generate()
parquet_files = sorted(Path(input_dir).glob("*.parquet"))
print(f"Found {len(parquet_files)} parquet files in {input_dir}")
frames = []
for f in parquet_files:
df = pd.read_parquet(f)
df["source_file"] = f.stem
frames.append(df)
combined = pd.concat(frames, ignore_index=True)
print(f"Combined: {len(combined)} rows, {combined.columns.tolist()}")
return combined
def filter_and_transform(df):
"""Apply filters and add computed columns."""
# Filter out health check endpoints
df = df[df["endpoint"] != "/api/health"].copy()
print(f"After filtering health checks: {len(df)} rows")
# Categorize status codes
df["status_category"] = pd.cut(
df["status_code"],
bins=[0, 199, 299, 399, 499, 599],
labels=["1xx", "2xx", "3xx", "4xx", "5xx"],
)
# Performance buckets
df["is_slow"] = (df["response_time_ms"] > 500).astype(int)
df["perf_bucket"] = pd.cut(
df["response_time_ms"],
bins=[0, 50, 200, 500, 1000, float("inf")],
labels=["fast", "normal", "slow", "very_slow", "timeout"],
)
# Extract hour from timestamp
df["hour"] = df["timestamp"].dt.hour
df["day_of_week"] = df["timestamp"].dt.dayofweek
return df
def compute_summaries(df):
"""Compute endpoint and region summaries."""
endpoint_summary = df.groupby("endpoint").agg(
request_count=("user_id", "count"),
unique_users=("user_id", "nunique"),
avg_response_ms=("response_time_ms", "mean"),
p95_response_ms=("response_time_ms", lambda x: x.quantile(0.95)),
error_count=("is_slow", "sum"),
total_bytes=("bytes_sent", "sum"),
).reset_index()
region_summary = df.groupby("region").agg(
request_count=("user_id", "count"),
avg_response_ms=("response_time_ms", "mean"),
cache_hit_rate=("is_cached", "mean"),
).reset_index()
print(f"Endpoint summary:\n{endpoint_summary.to_string(index=False)}")
print(f"\nRegion summary:\n{region_summary.to_string(index=False)}")
return endpoint_summary, region_summary
def write_partitioned(df, output_dir):
"""Write partitioned parquet output by region."""
output_path = Path(output_dir)
if output_path.exists():
import shutil
shutil.rmtree(output_path)
output_path.mkdir(parents=True)
# Convert categoricals to string for parquet compatibility
for col in df.select_dtypes(include=["category"]).columns:
df[col] = df[col].astype(str)
for region, group in df.groupby("region"):
region_dir = output_path / f"region={region}"
region_dir.mkdir(exist_ok=True)
out_file = region_dir / "data.parquet"
group.to_parquet(out_file, index=False)
print(f"Wrote {out_file} ({len(group)} rows)")
def write_summaries(endpoint_summary, region_summary, output_dir):
"""Write summary tables as parquet."""
output_path = Path(output_dir)
output_path.mkdir(parents=True, exist_ok=True)
endpoint_summary.to_parquet(output_path / "endpoint_summary.parquet", index=False)
region_summary.to_parquet(output_path / "region_summary.parquet", index=False)
print(f"Wrote summary parquets to {output_path}")
def main():
df = load_all_parquet("raw_logs")
df = filter_and_transform(df)
endpoint_summary, region_summary = compute_summaries(df)
write_partitioned(df, "processed_logs")
write_summaries(endpoint_summary, region_summary, "processed_logs/summaries")
# Verify round-trip by reading back
read_back = pd.read_parquet("processed_logs/summaries/endpoint_summary.parquet")
print(f"\nRound-trip verification: {len(read_back)} endpoint summary rows read back")
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate synthetic retail sales data for pivot/melt operations."""
import os
import numpy as np
import pandas as pd
SEED = 42
N_ROWS = 60_000
def generate():
if os.path.exists("retail_sales.csv"):
return
rng = np.random.default_rng(SEED)
stores = [f"Store_{i:02d}" for i in range(1, 16)]
products = ["Laptop", "Phone", "Tablet", "Headphones", "Charger",
"Case", "Cable", "Monitor", "Keyboard", "Mouse"]
quarters = ["Q1", "Q2", "Q3", "Q4"]
years = [2022, 2023, 2024]
channels = ["online", "in-store", "phone"]
df = pd.DataFrame({
"transaction_id": range(N_ROWS),
"store": rng.choice(stores, N_ROWS),
"product": rng.choice(products, N_ROWS),
"year": rng.choice(years, N_ROWS),
"quarter": rng.choice(quarters, N_ROWS),
"channel": rng.choice(channels, N_ROWS, p=[0.5, 0.35, 0.15]),
"units_sold": rng.integers(1, 20, N_ROWS),
"revenue": np.round(rng.uniform(10, 2000, N_ROWS), 2),
"cost": np.round(rng.uniform(5, 1500, N_ROWS), 2),
"customer_satisfaction": rng.integers(1, 6, N_ROWS),
})
df["profit"] = df["revenue"] - df["cost"]
df.to_csv("retail_sales.csv", index=False)
print(f"Generated {len(df)} retail sales rows -> retail_sales.csv")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Pivot, melt, stack/unstack, and cross-tabulation on retail data.
Demonstrates various DataFrame reshape operations for sales analysis.
"""
import numpy as np
import pandas as pd
from generate_data import generate
def load_data():
generate()
df = pd.read_csv("retail_sales.csv")
print(f"Loaded {len(df)} retail sales rows")
return df
def pivot_revenue_by_product_quarter(df):
"""Pivot table: average revenue by product and quarter."""
pivot = pd.pivot_table(
df,
values="revenue",
index="product",
columns="quarter",
aggfunc="mean",
fill_value=0,
)
pivot = pivot.round(2)
print(f"Revenue pivot (product x quarter):\n{pivot}")
return pivot
def pivot_multi_agg(df):
"""Pivot table with multiple aggregation functions."""
pivot = pd.pivot_table(
df,
values=["revenue", "units_sold"],
index=["store"],
columns=["year"],
aggfunc={"revenue": ["sum", "mean"], "units_sold": "sum"},
fill_value=0,
)
print(f"Multi-agg pivot shape: {pivot.shape}")
print(f"Columns: {pivot.columns.tolist()[:8]}...")
return pivot
def melt_pivot_back(pivot_df):
"""Melt a pivoted DataFrame back to long format."""
# Reset index to make product a column
flat = pivot_df.reset_index()
melted = pd.melt(
flat,
id_vars=["product"],
var_name="quarter",
value_name="avg_revenue",
)
melted = melted.sort_values(["product", "quarter"])
print(f"Melted back to long format: {len(melted)} rows")
return melted
def stack_unstack_demo(df):
"""Demonstrate stack and unstack operations."""
# Create a multi-index aggregation
agg = df.groupby(["store", "product"]).agg(
total_revenue=("revenue", "sum"),
total_units=("units_sold", "sum"),
)
# Unstack product to columns
unstacked = agg["total_revenue"].unstack(fill_value=0)
print(f"Unstacked shape: {unstacked.shape}")
# Stack it back
stacked = unstacked.stack()
stacked.name = "total_revenue"
stacked = stacked.reset_index()
print(f"Re-stacked: {len(stacked)} rows")
return unstacked, stacked
def crosstab_analysis(df):
"""Cross-tabulation of channel vs product."""
# Count cross-tab
ct_count = pd.crosstab(
df["channel"],
df["product"],
margins=True,
margins_name="Total",
)
print(f"Count crosstab:\n{ct_count}")
# Value cross-tab (average satisfaction)
ct_sat = pd.crosstab(
df["channel"],
df["product"],
values=df["customer_satisfaction"],
aggfunc="mean",
).round(2)
print(f"\nSatisfaction crosstab:\n{ct_sat}")
# Normalized cross-tab
ct_norm = pd.crosstab(
df["channel"],
df["product"],
normalize="index",
).round(4)
print(f"\nNormalized crosstab:\n{ct_norm}")
return ct_count, ct_sat, ct_norm
def year_over_year_pivot(df):
"""Pivot to compare year-over-year performance by store."""
yearly = df.groupby(["store", "year"]).agg(
revenue=("revenue", "sum"),
units=("units_sold", "sum"),
avg_profit=("profit", "mean"),
).reset_index()
# Pivot years to columns for side-by-side comparison
yoy = yearly.pivot_table(
index="store",
columns="year",
values="revenue",
aggfunc="sum",
fill_value=0,
)
yoy.columns = [f"revenue_{y}" for y in yoy.columns]
yoy = yoy.reset_index()
# Compute growth rates
if "revenue_2023" in yoy.columns and "revenue_2022" in yoy.columns:
yoy["growth_22_23"] = (
(yoy["revenue_2023"] - yoy["revenue_2022"]) / yoy["revenue_2022"]
).round(4)
if "revenue_2024" in yoy.columns and "revenue_2023" in yoy.columns:
yoy["growth_23_24"] = (
(yoy["revenue_2024"] - yoy["revenue_2023"]) / yoy["revenue_2023"]
).round(4)
print(f"Year-over-year:\n{yoy.head().to_string(index=False)}")
return yoy
def main():
df = load_data()
# Pivot operations
revenue_pivot = pivot_revenue_by_product_quarter(df)
multi_pivot = pivot_multi_agg(df)
# Melt
melted = melt_pivot_back(revenue_pivot)
# Stack / Unstack
unstacked, stacked = stack_unstack_demo(df)
# Cross-tabulation
ct_count, ct_sat, ct_norm = crosstab_analysis(df)
# Year-over-year pivot
yoy = year_over_year_pivot(df)
print(f"\nAll reshape operations completed successfully.")
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Text cleaning pipeline using pandas string operations.
Reads messy contact data and applies a series of string transformations:
lowercase, strip whitespace, regex extraction, contains checks, and replacements.
"""
import pandas as pd
from generate_data import generate
def load_data():
generate()
df = pd.read_csv("raw_contacts.csv")
df["notes"] = df["notes"].fillna("")
print(f"Loaded {len(df)} raw contacts")
return df
def clean_names(df):
"""Normalize first and last names."""
df["first_name"] = df["first_name"].str.strip().str.lower().str.title()
df["last_name"] = df["last_name"].str.strip().str.lower().str.title()
df["full_name"] = df["first_name"] + " " + df["last_name"]
return df
def clean_emails(df):
"""Strip and lowercase emails, extract domain."""
df["email"] = df["email"].str.strip().str.lower()
df["email_domain"] = df["email"].str.extract(r"@([a-z0-9\.\-]+)$", expand=False)
df["is_company_email"] = df["email_domain"].str.contains(
r"\.(org|net)$", regex=True
).astype(int)
return df
def normalize_phones(df):
"""Extract digits from phone numbers into a standard 10-digit format."""
digits = df["phone"].str.replace(r"[^\d]", "", regex=True)
# Remove leading '1' for 11-digit US numbers
digits = digits.str.replace(r"^1(\d{10})$", r"\1", regex=True)
df["phone_clean"] = (
"(" + digits.str[:3] + ") " + digits.str[3:6] + "-" + digits.str[6:10]
)
return df
def parse_addresses(df):
"""Extract state and zip from address strings."""
df["address"] = df["address"].str.strip()
df["state"] = df["address"].str.extract(r",\s*([A-Z]{2})\s+\d{5}", expand=False)
df["zipcode"] = df["address"].str.extract(r"(\d{5})\s*$", expand=False)
return df
def process_notes(df):
"""Extract reference numbers, detect flags, clean up notes."""
df["notes"] = df["notes"].str.strip()
# Extract reference numbers like Ref#12345 or REF#99887
df["ref_number"] = df["notes"].str.extract(
r"[Rr][Ee][Ff]#(\d+)", expand=False
)
# Flag rows
df["is_vip"] = df["notes"].str.contains("VIP", case=False, na=False).astype(int)
df["has_bounced"] = df["notes"].str.contains("BOUNCED", case=False, na=False).astype(int)
df["needs_followup"] = df["notes"].str.contains(
"follow-up|pending", case=False, regex=True, na=False
).astype(int)
# Redact discount details
df["notes_redacted"] = df["notes"].str.replace(
r"Discount:\s*\d+%", "Discount: [REDACTED]", regex=True
)
return df
def summarize(df):
"""Print summary statistics about the cleaned data."""
print(f"\nCleaned {len(df)} contacts")
print(f" Unique domains: {df['email_domain'].nunique()}")
print(f" Company emails: {df['is_company_email'].sum()}")
print(f" VIP customers: {df['is_vip'].sum()}")
print(f" Bounced emails: {df['has_bounced'].sum()}")
print(f" With ref numbers: {df['ref_number'].notna().sum()}")
print(f" States found: {df['state'].nunique()}")
def main():
df = load_data()
df = clean_names(df)
df = clean_emails(df)
df = normalize_phones(df)
df = parse_addresses(df)
df = process_notes(df)
summarize(df)
df.to_csv("cleaned_contacts.csv", index=False)
print("\nWrote cleaned_contacts.csv")
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate synthetic messy text data for string operations."""
import os
import numpy as np
import pandas as pd
SEED = 42
N_ROWS = 30_000
def generate():
if os.path.exists("raw_contacts.csv"):
return
rng = np.random.default_rng(SEED)
first_names = ["Alice", "Bob", " Charlie", "Diana ", " Eve", "FRANK",
"grace", " HANK ", "Ivy", " jack"]
last_names = ["Smith", " JONES", "Williams ", " BROWN", "davis",
" Miller", "WILSON ", "moore", " Taylor", "Anderson"]
domains = ["gmail.com", "yahoo.com", "outlook.com", "company.org", "example.net"]
phones_raw = []
emails_raw = []
addresses_raw = []
for _ in range(N_ROWS):
# messy phone: mix of formats
area = rng.integers(200, 999)
mid = rng.integers(100, 999)
last4 = rng.integers(1000, 9999)
fmt = rng.choice(["paren", "dash", "dot", "plain", "intl"])
if fmt == "paren":
phones_raw.append(f"({area}) {mid}-{last4}")
elif fmt == "dash":
phones_raw.append(f"{area}-{mid}-{last4}")
elif fmt == "dot":
phones_raw.append(f"{area}.{mid}.{last4}")
elif fmt == "plain":
phones_raw.append(f"{area}{mid}{last4}")
else:
phones_raw.append(f"+1-{area}-{mid}-{last4}")
fn = rng.choice(first_names)
ln = rng.choice(last_names)
dom = rng.choice(domains)
emails_raw.append(f" {fn.strip().lower()}.{ln.strip().lower()}@{dom} ")
num = rng.integers(1, 9999)
street = rng.choice(["Main St", "Oak Ave", "1st Blvd", "Elm Dr", "Pine Ln"])
state = rng.choice(["CA", "NY", "TX", "FL", "WA", "IL"])
zipcode = rng.integers(10000, 99999)
addresses_raw.append(f" {num} {street}, {state} {zipcode} ")
df = pd.DataFrame({
"first_name": rng.choice(first_names, N_ROWS),
"last_name": rng.choice(last_names, N_ROWS),
"email": emails_raw,
"phone": phones_raw,
"address": addresses_raw,
"notes": rng.choice([
"VIP customer - priority support",
"CALLED 2024-01-15: billing issue",
"Ref#12345 - pending review",
" no notes ",
"email BOUNCED on 2024-03-01",
"Discount: 20% off next order",
"REF#99887 follow-up required",
"",
], N_ROWS),
})
df.to_csv("raw_contacts.csv", index=False)
print(f"Generated {len(df)} messy contact rows -> raw_contacts.csv")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate synthetic sensor data with minute-level timestamps."""
import os
import numpy as np
import pandas as pd
SEED = 42
N_MINUTES = 60_000 # ~41 days of minute-level data
def generate():
if os.path.exists("sensor_data.csv"):
return
rng = np.random.default_rng(SEED)
timestamps = pd.date_range(
start="2024-01-01", periods=N_MINUTES, freq="min"
)
# Simulate three sensors with seasonal patterns and noise
hour_of_day = timestamps.hour + timestamps.minute / 60.0
day_cycle = np.sin(2 * np.pi * hour_of_day / 24.0)
df = pd.DataFrame({
"timestamp": timestamps,
"sensor_id": rng.choice(["S1", "S2", "S3"], N_MINUTES),
"temperature": 20.0 + 5.0 * day_cycle + rng.normal(0, 0.5, N_MINUTES),
"humidity": 60.0 - 10.0 * day_cycle + rng.normal(0, 2.0, N_MINUTES),
"pressure": 1013.0 + rng.normal(0, 3.0, N_MINUTES),
"voltage": 3.3 + rng.normal(0, 0.05, N_MINUTES),
})
df["temperature"] = np.round(df["temperature"], 2)
df["humidity"] = np.clip(np.round(df["humidity"], 1), 0, 100)
df["pressure"] = np.round(df["pressure"], 1)
df["voltage"] = np.round(df["voltage"], 3)
df.to_csv("sensor_data.csv", index=False)
print(f"Generated {len(df)} sensor readings -> sensor_data.csv")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Timeseries resampling and rolling statistics pipeline.
Reads minute-level sensor data, resamples to hourly and daily frequencies,
and computes rolling window statistics for anomaly detection thresholds.
"""
import numpy as np
import pandas as pd
from generate_data import generate
def load_data():
generate()
df = pd.read_csv("sensor_data.csv", parse_dates=["timestamp"])
df = df.sort_values("timestamp")
print(f"Loaded {len(df)} sensor readings from "
f"{df['timestamp'].min()} to {df['timestamp'].max()}")
return df
def resample_hourly(df):
"""Resample each sensor to hourly frequency."""
hourly_frames = []
for sensor_id, group in df.groupby("sensor_id"):
ts = group.set_index("timestamp")
hourly = ts[["temperature", "humidity", "pressure", "voltage"]].resample("h").agg(
["mean", "min", "max", "std"]
)
# Flatten multi-level columns
hourly.columns = ["_".join(col) for col in hourly.columns]
hourly["sensor_id"] = sensor_id
hourly = hourly.reset_index()
hourly_frames.append(hourly)
result = pd.concat(hourly_frames, ignore_index=True)
print(f"Hourly resampled: {len(result)} rows")
return result
def resample_daily(df):
"""Resample all sensors to daily frequency with aggregation."""
ts = df.set_index("timestamp")
daily = ts.groupby("sensor_id").resample("D").agg({
"temperature": ["mean", "min", "max"],
"humidity": ["mean", "min", "max"],
"pressure": "mean",
"voltage": "mean",
})
daily.columns = ["_".join(col) for col in daily.columns]
daily = daily.reset_index()
print(f"Daily resampled: {len(daily)} rows")
return daily
def compute_rolling_stats(hourly):
"""Compute rolling 24-hour statistics on the hourly data."""
rolling_frames = []
for sensor_id, group in hourly.groupby("sensor_id"):
g = group.sort_values("timestamp").copy()
g["temp_rolling_mean_24h"] = (
g["temperature_mean"].rolling(window=24, min_periods=6).mean()
)
g["temp_rolling_std_24h"] = (
g["temperature_mean"].rolling(window=24, min_periods=6).std()
)
g["humidity_rolling_mean_24h"] = (
g["humidity_mean"].rolling(window=24, min_periods=6).mean()
)
g["pressure_expanding_mean"] = g["pressure_mean"].expanding(min_periods=1).mean()
# Anomaly flag: temperature deviates more than 2 std from rolling mean
g["temp_anomaly"] = (
(g["temperature_mean"] - g["temp_rolling_mean_24h"]).abs()
> 2 * g["temp_rolling_std_24h"]
).astype(int)
rolling_frames.append(g)
result = pd.concat(rolling_frames, ignore_index=True)
anomaly_count = result["temp_anomaly"].sum()
print(f"Rolling stats computed; {anomaly_count} temperature anomalies detected")
return result
def compute_daily_change(daily):
"""Compute day-over-day changes using shift."""
change_frames = []
for sensor_id, group in daily.groupby("sensor_id"):
g = group.sort_values("timestamp").copy()
g["temp_change"] = g["temperature_mean"] - g["temperature_mean"].shift(1)
g["humidity_change"] = g["humidity_mean"] - g["humidity_mean"].shift(1)
g["temp_cummax"] = g["temperature_max"].cummax()
g["temp_cummin"] = g["temperature_min"].cummin()
change_frames.append(g)
result = pd.concat(change_frames, ignore_index=True)
print(f"Daily changes computed for {result['sensor_id'].nunique()} sensors")
return result
def main():
df = load_data()
hourly = resample_hourly(df)
daily = resample_daily(df)
hourly_with_rolling = compute_rolling_stats(hourly)
daily_with_changes = compute_daily_change(daily)
print(f"\nHourly sample:\n{hourly_with_rolling.head(3).to_string(index=False)}")
print(f"\nDaily sample:\n{daily_with_changes.head(3).to_string(index=False)}")
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Generate synthetic stock trading data for window function analysis."""
import os
import numpy as np
import pandas as pd
SEED = 42
N_DAYS = 500
N_STOCKS = 50
def generate():
if os.path.exists("stock_trades.csv"):
return
rng = np.random.default_rng(SEED)
dates = pd.bdate_range(start="2022-01-03", periods=N_DAYS)
tickers = [f"STK{i:03d}" for i in range(N_STOCKS)]
rows = []
for ticker in tickers:
base_price = rng.uniform(10, 500)
prices = [base_price]
for _ in range(N_DAYS - 1):
change = rng.normal(0, base_price * 0.02)
prices.append(max(1.0, prices[-1] + change))
for i, date in enumerate(dates):
rows.append({
"date": date,
"ticker": ticker,
"close": round(prices[i], 2),
"volume": int(rng.integers(10_000, 5_000_000)),
"high": round(prices[i] * (1 + rng.uniform(0, 0.03)), 2),
"low": round(prices[i] * (1 - rng.uniform(0, 0.03)), 2),
})
df = pd.DataFrame(rows)
df["trade_value"] = df["close"] * df["volume"]
df.to_csv("stock_trades.csv", index=False)
print(f"Generated {len(df)} stock trade rows -> stock_trades.csv")
if __name__ == "__main__":
generate()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Window function analysis on stock trading data.
Computes rankings, cumulative sums, rolling averages, expanding statistics,
and shift/lag features for each stock ticker.
"""
import numpy as np
import pandas as pd
from generate_data import generate
def load_data():
generate()
df = pd.read_csv("stock_trades.csv", parse_dates=["date"])
df = df.sort_values(["ticker", "date"])
print(f"Loaded {len(df)} trades for {df['ticker'].nunique()} tickers")
return df
def add_rankings(df):
"""Rank stocks by close price and volume within each date."""
df["price_rank_dense"] = df.groupby("date")["close"].rank(
method="dense", ascending=False
)
df["price_rank_min"] = df.groupby("date")["close"].rank(
method="min", ascending=False
)
df["volume_rank"] = df.groupby("date")["volume"].rank(
method="average", ascending=False
)
df["price_pctrank"] = df.groupby("date")["close"].rank(pct=True)
print(f"Rankings added; top stock on last day: "
f"rank 1 = {df.loc[df['price_rank_dense'] == 1].tail(1)['ticker'].values}")
return df
def add_cumulative(df):
"""Compute cumulative statistics per ticker."""
df["cumsum_volume"] = df.groupby("ticker")["volume"].cumsum()
df["cumsum_trade_value"] = df.groupby("ticker")["trade_value"].cumsum()
df["cummax_close"] = df.groupby("ticker")["close"].cummax()
df["cummin_close"] = df.groupby("ticker")["close"].cummin()
df["cum_avg_close"] = df["cumsum_trade_value"] / df["cumsum_volume"]
print("Cumulative stats added")
return df
def add_rolling_stats(df):
"""Compute rolling window statistics per ticker."""
rolling_frames = []
for ticker, group in df.groupby("ticker"):
g = group.sort_values("date").copy()
# 5-day and 20-day rolling averages
g["sma_5"] = g["close"].rolling(window=5, min_periods=1).mean()
g["sma_20"] = g["close"].rolling(window=20, min_periods=5).mean()
# Rolling standard deviation (volatility)
g["volatility_20"] = g["close"].rolling(window=20, min_periods=5).std()
# Rolling min/max (support/resistance levels)
g["rolling_high_20"] = g["high"].rolling(window=20, min_periods=5).max()
g["rolling_low_20"] = g["low"].rolling(window=20, min_periods=5).min()
# Rolling sum of volume
g["volume_sum_10"] = g["volume"].rolling(window=10, min_periods=1).sum()
rolling_frames.append(g)
result = pd.concat(rolling_frames, ignore_index=True)
print("Rolling stats added (SMA-5, SMA-20, volatility, support/resistance)")
return result
def add_expanding_stats(df):
"""Compute expanding window statistics per ticker."""
expanding_frames = []
for ticker, group in df.groupby("ticker"):
g = group.sort_values("date").copy()
g["expanding_mean"] = g["close"].expanding(min_periods=1).mean()
g["expanding_std"] = g["close"].expanding(min_periods=2).std()
g["expanding_max"] = g["close"].expanding(min_periods=1).max()
g["expanding_min"] = g["close"].expanding(min_periods=1).min()
expanding_frames.append(g)
result = pd.concat(expanding_frames, ignore_index=True)
print("Expanding stats added")
return result
def add_shift_features(df):
"""Compute lag/lead features and returns."""
shift_frames = []
for ticker, group in df.groupby("ticker"):
g = group.sort_values("date").copy()
# Lag features
g["prev_close"] = g["close"].shift(1)
g["prev_close_5"] = g["close"].shift(5)
# Daily return
g["daily_return"] = (g["close"] - g["prev_close"]) / g["prev_close"]
# 5-day return
g["return_5d"] = (g["close"] - g["prev_close_5"]) / g["prev_close_5"]
# Lead (next day close)
g["next_close"] = g["close"].shift(-1)
# Diff
g["close_diff"] = g["close"].diff()
g["volume_diff"] = g["volume"].diff()
shift_frames.append(g)
result = pd.concat(shift_frames, ignore_index=True)
print("Shift/lag features added (returns, diffs, leads)")
return result
def generate_signals(df):
"""Simple moving average crossover signals."""
df["sma_cross"] = (df["sma_5"] > df["sma_20"]).astype(int)
df["signal_change"] = df.groupby("ticker")["sma_cross"].diff().fillna(0).astype(int)
buy_signals = (df["signal_change"] == 1).sum()
sell_signals = (df["signal_change"] == -1).sum()
print(f"Signals: {buy_signals} buys, {sell_signals} sells")
return df
def main():
df = load_data()
df = add_rankings(df)
df = add_cumulative(df)
df = add_rolling_stats(df)
df = add_expanding_stats(df)
df = add_shift_features(df)
df = generate_signals(df)
print(f"\nFinal shape: {df.shape}")
sample = df[df["ticker"] == "STK000"].tail(5)
print(f"\nSample (STK000 last 5 days):\n"
f"{sample[['date', 'close', 'sma_5', 'sma_20', 'daily_return', 'price_rank_dense']].to_string(index=False)}")
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Small PyTorch training script."""
from __future__ import annotations
import torch
from torch import nn
from torch.utils.data import DataLoader, TensorDataset
class Model(nn.Module):
def __init__(self) -> None:
super().__init__()
self.net = nn.Sequential(
nn.Linear(1024, 4096),
nn.ReLU(),
nn.Linear(4096, 4096),
nn.ReLU(),
nn.Linear(4096, 10),
)
def forward(self, x: torch.Tensor) -> torch.Tensor:
return self.net(x)
def main() -> None:
device = "cuda" if torch.cuda.is_available() else "cpu"
x = torch.randn(20_000, 1024)
y = torch.randint(0, 10, (20_000,))
loader = DataLoader(TensorDataset(x, y), batch_size=64, shuffle=True, num_workers=0)
model = Model().to(device)
opt = torch.optim.AdamW(model.parameters(), lr=1e-3)
loss_fn = nn.CrossEntropyLoss()
for epoch in range(2):
for xb, yb in loader:
xb = xb.to(device)
yb = yb.to(device)
opt.zero_grad(set_to_none=True)
loss = loss_fn(model(xb), yb)
loss.backward()
opt.step()
print(f"epoch={epoch} loss={float(loss.detach().cpu()):.4f}")
if __name__ == "__main__":
main()
# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
"""Pandas nullable-value cleanup pipeline."""
from __future__ import annotations
import numpy as np
import pandas as pd
def build_frame() -> pd.DataFrame:
return pd.DataFrame(
{
"account": pd.Series([1, 2, None, 4, 5, None], dtype="Int64"),
"region": pd.Series(["west", None, "east", "west", None, "east"], dtype="string"),
"score": [0.8, np.nan, 0.3, 0.9, np.nan, 0.2],
"tier": pd.Series(["gold", "silver", None, "gold", "bronze", None], dtype="string"),
}
)
def clean(frame: pd.DataFrame) -> pd.DataFrame:
result = frame.copy()
result["region"] = result["region"].fillna("unknown")
result["tier"] = result["tier"].fillna("unassigned")
result["score"] = result["score"].where(result["score"].notna(), result["score"].median())
result["high_score"] = result["score"] >= 0.75
grouped = (
result.groupby("region", dropna=False)
.agg(
accounts=("account", "count"),
avg_score=("score", "mean"),
high_count=("high_score", "sum"),
)
.reset_index()
.sort_values("region")
)
return grouped
def main() -> None:
print(clean(build_frame()).to_string(index=False))
if __name__ == "__main__":
main()
Attribution
This task is source-inspired by cuDF null-handling tests.
- Source: https://github.com/rapidsai/cudf/blob/235f69a6fcef/python/cudf/cudf/tests/dataframe/methods/test_fillna.py
- Upstream project: RAPIDS cuDF
- License: Apache-2.0
- Local changes: original pandas fixture written for benchmark scoring; no upstream code copied.
cuDF API Patterns, Gaps, and Semantic Differences
Key Semantic Differences from pandas
Null/NaN Handling
cuDF preserves nullable dtypes more often than pandas and uses Arrow-style nulls instead of float NaN promotion for nullable numeric columns:
import cudf
import pandas as pd
s = cudf.Series([1, None, 3])
print(s.dtype) # Int64 (nullable), not float64 with NaN
# Check for null
s.isnull() # works as expected
s.isna() # equivalent
# Fill nulls
s.fillna(0) # worksDifference: pd.Series([1, None, 3]) → dtype float64 with NaN; cuDF → nullable Int64 with <NA>.
For string columns in current releases, missing string values display as None rather than <NA>. Do not write tests that depend on the display repr; compare with .isna(), .notna(), or typed result values.
When comparing cuDF output with a pandas nullable reference, convert with nullable=True:
actual_pdf = gdf.to_pandas(nullable=True)This keeps nullable pandas dtypes when they exist, instead of converting nulls to np.nan or None during the comparison boundary.
For a null-heavy workflow, keep the pandas behavior as a compact reference and make the GPU path explicit:
- scalar, dictionary, forward, and backward fills map directly to cuDF
- group-specific fills are usually
groupby().transform(...)followed by
fillna(...)
- conditional fills are boolean masks plus assignment, or a grouped aggregate
merged back onto the original frame
- linear interpolation is a semantic boundary; use cuDF only after checking the
installed API behavior, or keep that narrow step under cudf.pandas with a parity check
Validate row count, null count by column, representative filled values, grouped aggregates, and any rows produced by sort/interpolation-sensitive code.
Sort Stability
cuDF sort is not stable by default:
# Unstable (default) — faster
df.sort_values("col")
# Stable — required when sort order must match pandas exactly
df.sort_values("col", stable=True)String Operations — RE2 Regex
cuDF uses RE2 (not Python's re / PCRE). Some patterns differ:
# RE2 does not support:
# - Lookahead/lookbehind: (?=...), (?!...)
# - Backreferences: \1
# - Possessive quantifiers: ?+, *+
# RE2-compatible (works):
df["col"].str.contains(r"\d+")
df["col"].str.replace(r"[aeiou]", "", regex=True)
# Not RE2-compatible (will fail or fall back):
df["col"].str.contains(r"(?=.*foo)") # lookahead — use different approachCuPy Array Output
When you access .values on a cuDF Series/DataFrame, you get a CuPy array (not NumPy):
import cudf
import cupy as cp
df = cudf.DataFrame({"a": [1, 2, 3]})
arr = df["a"].values # CuPy array, not NumPy!
type(arr) # <class 'cupy.ndarray'>
# To get NumPy explicitly:
np_arr = df["a"].to_numpy()
np_arr = cp.asnumpy(arr)Common API Gaps and Workarounds
The pandas API surface is vast and cuDF only covers a limited subset of it. This section lays out some of the common gaps but it should not be construed as an exhaustive list of discrepancies between the cuDF and pandas APIs.
Operations Not Yet in cuDF
| pandas Operation | Status | Workaround |
|---|---|---|
df.apply(func, axis=0) | Column-wise apply: limited | Rewrite as vectorized cuDF ops |
df.apply(func, axis=1) | Row-wise apply: limited | Use df.apply() for simple funcs; otherwise cudf.pandas fallback |
Some pd.Grouper options | Partial | Use resample or direct groupby |
pd.read_html() | Not supported | Use pandas, then cudf.from_pandas() |
pd.ExcelWriter / read_excel | Not supported | Convert to CSV/Parquet first |
df.to_sql() | Not supported | Convert to pandas, then use pandas |
| Multi-level columns (MultiIndex) | Partial | Flatten column names first |
Reshape and Crosstab Fidelity
cudf.pivot_table, cudf.melt, cudf.crosstab, DataFrame.unstack, and DataFrame.stack cover many reshape workflows. Treat the source pandas schema as observable behavior when a pipeline depends on reshape output:
- Capture expected index labels, column labels or levels, names, shape, and
representative values from the pandas path before rewriting.
- Preserve pandas MultiIndex columns when the downstream code consumes them. If
a flat schema is the practical cuDF representation, return a documented mapping such as revenue_sum_2024 and validate consumers against that schema.
- For multi-aggregation
pivot_tableoutputs, keep aggregation names in the
schema. Build the cuDF result from explicit grouped aggregations when needed, then either recreate the pandas column levels or flatten with deterministic names such as {value}_{agg}_{column}.
- Implement missing
crosstabconveniences with explicit GPU operations:
counts via cudf.crosstab, margins via row/column sums, and row-normalized values by dividing each row by its row total.
- Use
cudf.pandasas a compatibility-first path when exact pandas reshape
semantics are the goal and explicit cuDF would require broad schema changes.
- Add a reusable validation helper that compares shape, index/column labels,
aggregation names, null placement, and numeric values against the pandas reference on a small fixture.
Time-Series and Rolling Fidelity
cuDF supports datetime columns, sorting, grouped operations, shifts, cumulative operations, and many rolling-window patterns. Preserve pandas-visible time semantics when rewriting:
- keep timezone, timestamp dtype, frequency, and bucket labels as part of the
output contract
- sort by grouping keys and timestamp before grouped
shift,rolling,
cumulative, or expanding-style calculations
- validate sparse or missing buckets against the pandas reference; explicitly
materialize the desired bucket grid when downstream consumers expect empty periods
- use final
.to_pandas()only for display, plotting, or reference comparison
I/O Formats Supported by cuDF
# Fully supported (fast GPU I/O)
cudf.read_csv(), cudf.read_parquet(), cudf.read_json()
cudf.read_orc(), cudf.read_feather(), cudf.read_avro()
# Not supported (use pandas, convert with cudf.from_pandas())
# Excel, HTML, SQL, HDF5, SAS, Stata, pickleUseful cuDF-Specific APIs
# Convert between pandas and cuDF
cudf_df = cudf.from_pandas(pd_df)
pd_df = cudf_df.to_pandas()
# Interop with CuPy
import cupy as cp
arr = cp.asarray(df["col"]) # zero-copy view
df["new_col"] = cudf.Series(arr) # back to cuDFPerformance Tips
1. Cast to float32 early: df[numeric_cols] = df[numeric_cols].astype("float32") 2. Use `cudf.read_parquet()` not CSV: Parquet is columnar and dramatically faster to read 3. Avoid `.apply()` with Python lambdas: Use built-in cuDF ops instead 4. Use `persist()` with dask-cuDF: keeps computed data on GPU workers to avoid recomputation 5. Avoid mid-pipeline `.to_pandas()`: each roundtrip is a PCIe transfer
cudf.pandas Accelerator — Deep Dive
How It Works
cudf.pandas replaces the pandas module with a proxy that routes operations to cuDF when supported, falling back to standard pandas on CPU silently for unsupported operations. The fallback is transparent — code continues to work correctly, but unsupported ops run on CPU.
Activation Methods
| Method | Use Case |
|---|---|
%load_ext cudf.pandas | Jupyter/IPython notebooks |
python -m cudf.pandas script.py | CLI script execution |
import cudf.pandas; cudf.pandas.install() | Programmatic, multiprocessing |
Critical: Activation must happen BEFORE any pandas import, direct or transitive. If you're using IPython and pandas was already imported in the kernel, restart and run activation first. Direct usage of cudf.pandas.install() in a script cannot be undone and the script must be restarted.
Profiling for GPU vs CPU Ops
Cell-Level Profiling (Jupyter)
%load_ext cudf.pandas
import pandas as pd
%%cudf.pandas.profile
df = pd.read_csv("data.csv")
result = df.groupby("category")["amount"].sum()
df.merge(lookup, on="id")Output shows each operation's execution path (GPU or CPU) and time.
Line-Level Profiling
%%cudf.pandas.line_profile
df = pd.DataFrame({"a": range(1000000), "b": range(1000000)})
result = df.groupby("a")["b"].sum() # shows GPU time
df.apply(lambda x: x + 1, axis=1) # shows CPU fallback timeCLI Profiling
python -m cudf.pandas --profile my_script.pyDetecting Silent Fallbacks
The profiling tools are also a convenient way to detect silent fallback. If the profiles show tasks running on the CPU unexpectedly, you may be hitting unsupported GPU methods (limitations are discussed in depth in the api-patterns.md reference file). Try reproducing with raw cudf code without cudf.pandas to verify.
Verifying GPU Is Actually Used
# Method 1: Run nvidia-smi during execution
# nvidia-smi dmon -s u -d 1
# Method 2: Check cudf.pandas stats
import cudf.pandas
stats = cudf.pandas.get_stats()
print(stats) # shows GPU vs CPU operation countsIf GPU utilization stays 0% during execution, the entire workload fell back. Diagnose with %%cudf.pandas.profile.
multiprocessing Support
# This pattern ensures workers also use cudf.pandas
import cudf.pandas
cudf.pandas.install() # must be FIRST, before everything else
from multiprocessing import Pool
import pandas as pd
def process_chunk(args):
# Workers inherit cudf.pandas installation
df = pd.read_csv(args)
return df.groupby("key")["value"].sum()
with Pool(4) as pool:
results = pool.map(process_chunk, file_list)Limitations
- Usage of the NumPy C API: Many projects have custom extension modules that interface with pandas dataframes via the NumPy C API for interacting with individual pandas columns. That will never work with cudf.pandas.
- Subclassed DataFrames: code that subclasses
pd.DataFramemay not work with cudf.pandas proxy - Private pandas APIs (
pd._libs.*, etc.): not supported - In-place operations with external code: if third-party code holds references to pandas internals, proxy may not intercept correctly
- cudf.pandas does not speed up Python-level loops: vectorize first, then accelerate
When to Move to Explicit cuDF
Move from cudf.pandas to explicit cuDF when: 1. Profile shows >30% CPU fallback rate on hot paths 2. You need cuDF-specific features (e.g., cudf.set_option("spill", True)) 3. You need explicit control over dtype casting (float32 optimization) 4. You're building a cuDF-first library, not accelerating existing pandas code
dask-cuDF Patterns
Preferred API: dask.dataframe Backend (release 24.06+)
The recommended way to use dask-cuDF is via the dask.dataframe backend config, not import dask_cudf directly. The backend API enables the query planning optimizer (predicate pushdown, projection pushdown) introduced in release 24.06+.
import dask
dask.config.set({"dataframe.backend": "cudf"})
import dask.dataframe as dd
# Read — now GPU-backed with query planning
ddf = dd.read_parquet("data/*.parquet")
ddf = dd.read_csv("data/*.csv")
# All standard dask.dataframe operations work
result = ddf.groupby("key")["value"].sum()Explicit `dask_cudf` import is still valid but bypasses query planning:
import dask_cudf # works, but no optimizer — use for legacy code only
ddf = dask_cudf.read_parquet("data/*.parquet")Cluster Setup
Always use LocalCUDACluster, even for a single GPU — it pins GPU affinity, enables the dashboard, and is required for proper spill configuration:
from dask_cuda import LocalCUDACluster
from dask.distributed import Client
import dask
dask.config.set({"dataframe.backend": "cudf"})
# Standard setup — one worker per GPU
cluster = LocalCUDACluster(
enable_cudf_spill=True, # cuDF-native spill; preferred over device_memory_limit
rmm_pool_size=0.8, # leave headroom for non-RMM allocations
)
client = Client(cluster)
# With UCX automatic transport selection for communication-heavy workloads
cluster = LocalCUDACluster(
enable_cudf_spill=True,
rmm_pool_size=0.8,
protocol="ucx",
)Partition Sizing
Partition size is the most impactful tuning parameter:
| Workload | Target Partition Size |
|---|---|
| General ETL | 1/32 – 1/8 of single GPU memory |
| Shuffle-intensive (groupby, join, sort) | 1/32 – 1/16 of GPU memory |
# Check current partitions
print(f"Partitions: {ddf.npartitions}")
# Tune at read time (most efficient)
ddf = dd.read_parquet("data/", blocksize="256MB") # adjust to hit target partition size
# Repartition after load if needed
ddf = ddf.repartition(npartitions=64)Reading Data
Local Parquet (Recommended)
import dask.dataframe as dd
# Project only needed columns — pushed down to storage
ddf = dd.read_parquet("data/*.parquet", columns=["col1", "col2", "key"])
# aggregate_files=True merges small files into larger partitions
ddf = dd.read_parquet("data/", aggregate_files=True, blocksize="512MB")Remote Storage (S3, GCS)
# Use blocksize=None to avoid slow metadata collection on remote stores
ddf = dd.read_parquet(
"s3://bucket/prefix/",
blocksize=None,
filesystem="arrow", # pyarrow filesystem for S3/GCS
columns=["col1", "col2"],
)Aggregation Patterns
Low-cardinality groupby
# split_out=1 avoids unnecessary shuffle for few output groups
result = ddf.groupby("status_code")["amount"].sum(split_out=1)High-cardinality groupby (default)
result = ddf.groupby("customer_id").agg({"amount": "sum", "count": "count"})Join / Merge Patterns
# Standard join (both datasets distributed)
merged = large_ddf.merge(other_large_ddf, on="id", how="left")
# Small table join: broadcast=True avoids shuffling the large table
merged = large_ddf.merge(
small_lookup_df, # cuDF DataFrame or small dask-cuDF
on="id",
how="left",
broadcast=True, # sends small_lookup to all workers; no shuffle
)Sort vs. Shuffle
# sort_values is expensive — triggers full shuffle + materialization
# AVOID unless you actually need a globally ordered output:
sorted_ddf = ddf.sort_values("timestamp") # use sparingly
# If you need rows grouped by key (not sorted), use shuffle instead:
from dask_cudf import shuffle
shuffled = shuffle(ddf, on="customer_id") # redistributes by key, much cheaperBuilding Distributed Collections
# Preferred: from_map enables column projection pushdown
from dask.dataframe import from_map
import cudf
def load_partition(path, columns=None):
return cudf.read_parquet(path, columns=columns)
files = ["data/part_0000.parquet", "data/part_0001.parquet"]
ddf = from_map(
load_partition,
files,
meta=cudf.read_parquet(files[0], nrows=0), # avoids eager first-partition read
)
# from_delayed works but loses projection pushdown
from dask import delayed
parts = [delayed(cudf.read_parquet)(f) for f in files]
ddf = dask_cudf.from_delayed(parts) # fallback if from_map doesn't applyEager Execution Traps
These calls trigger immediate computation — avoid mid-pipeline:
| Call | Why it's expensive |
|---|---|
.compute() on large collection | Pulls all data to one GPU |
.persist() without client.wait() | Silent if client not set up |
len(ddf) | Full scan |
ddf.head() / ddf.tail() | Materializes first/last partition |
ddf.sort_values(...) | Full shuffle |
ddf.set_index(col) | Full shuffle + sort |
Persist pattern (when you query the same data multiple times):
ddf = ddf.persist()
client.wait(ddf) # block until all partitions are in GPU memory
result1 = ddf[ddf["a"] > 0].compute()
result2 = ddf[ddf["b"] > 0].compute() # fast — data already in memoryNever call `.compute()` on a collection larger than single-GPU memory — it will OOM. Instead write to Parquet and read back in pieces.
Writing Results
# Parquet (recommended — partitioned output)
ddf.to_parquet("output/", write_index=False)
# To single cuDF DataFrame — only when result fits in GPU memory
result_cudf = ddf.compute()
# To pandas — only at the very end for CPU or non-GPU handoff
result_pd = ddf.to_pandas()OOM Diagnosis
# Step 1: Check worker memory pressure from dashboard
print(client.dashboard_link) # open in browser → Workers tab
# Step 2: Increase partition count to reduce per-partition memory
ddf = ddf.repartition(npartitions=ddf.npartitions * 2)
# Step 3: If not already enabled, add cuDF-native spilling
# (restart cluster with enable_cudf_spill=True, rmm_pool_size=0.9)
# Step 4: Move filter/project before expensive operations
ddf = ddf[["needed_col1", "needed_col2", "key"]] # project first
ddf = ddf[ddf["amount"] > 0] # filter early
result = ddf.groupby("key")["needed_col1"].sum().compute()Anti-Patterns
For new dask-cuDF code, use the backend setup shown in the Preferred API section above. The examples here focus on execution and materialization mistakes after the backend has been selected.
# AVOID: calling .compute() mid-pipeline
intermediate = ddf.groupby("a")["b"].sum().compute() # breaks lazy graph
result = intermediate.groupby("c")["b"].mean() # now CPU pandas!
# CORRECT: chain lazily, compute once
result = (
ddf.groupby("a")["b"].sum()
.reset_index()
.groupby("c")["b"].mean()
.compute()
)
# AVOID: collecting huge dataset to display
print(ddf.compute()) # OOM risk
# CORRECT: sample or head
print(ddf.head(10)) # shows first 10 rows onlyDescription: <br>
Official NVIDIA-authored guidance for NVIDIA cuDF GPU DataFrames, pandas acceleration, dask-cuDF, ETL, joins, groupby, CSV/Parquet I/O, nullable semantics, and multi-GPU DataFrame workloads. <br>
This skill is ready for commercial/non-commercial use. <br>
Owner
NVIDIA <br>
License/Terms of Use: <br>
CC-BY-4.0 AND Apache-2.0 <br>
Use Case: <br>
Developers and engineers building GPU-accelerated data processing pipelines using NVIDIA cuDF for DataFrame operations, ETL, joins, groupby, CSV/Parquet I/O, nullable semantics, and multi-GPU workloads. <br>
Deployment Geography for Use: <br>
Global <br>
Known Risks and Mitigations: <br>
Risk: Review before execution as proposals could introduce incorrect or misleading guidance into skills. <br> Mitigation: Review and scan skill before deployment. <br>
Reference(s): <br>
- cuDF API Patterns, Gaps, and Semantic Differences <br>
- cudf.pandas Accelerator Deep Dive <br>
- dask-cuDF Patterns <br>
- cuDF Documentation <br>
- dask-cuDF API Reference <br>
- cuDF GitHub Repository <br>
Skill Output: <br>
Output Type(s): [Code, Configuration instructions] <br> Output Format: [Markdown with inline Python and bash code blocks] <br> Output Parameters: [1D] <br> Other Properties Related to Output: [None] <br>
Evaluation Agents Used: <br>
- claude-code <br>
- codex <br>
Evaluation Tasks: <br>
13 evaluation tasks (12 positive skill-activation, 1 negative) with 2 attempts per task; pass threshold 50%. <br>
Evaluation Metrics Used: <br>
Reported benchmark dimensions: <br>
- Security: Checks whether skill-assisted execution avoids unsafe behavior such as secret leakage, destructive commands, or unauthorized access. <br>
- Correctness: Checks whether the agent follows the expected workflow and produces the correct final output. <br>
- Discoverability: Checks whether the agent loads the skill when relevant and avoids using it when irrelevant. <br>
- Effectiveness: Checks whether the agent performs measurably better with the skill than without it. <br>
- Efficiency: Checks whether the agent uses fewer tokens and avoids redundant work. <br>
Underlying evaluation signals used in this run: <br>
security: Checks for unsafe operations, secret leakage, and unauthorized access. <br>skill_execution: Verifies that the agent loaded the expected skill and workflow. <br>skill_efficiency: Checks routing quality, decoy avoidance, and redundant tool usage. <br>accuracy: Grades final-answer correctness against the reference answer. <br>goal_accuracy: Checks whether the overall user task completed successfully. <br>behavior_check: Verifies expected behavior steps, including safety expectations. <br>token_efficiency: Compares token usage with and without the skill. <br>
Evaluation Results: <br>
| Dimension | Num | claude-code | codex |
|---|---|---|---|
| Security | 8 | 92% (+12%) | 100% (+0%) |
| Correctness | 8 | 96% (+10%) | 92% (+8%) |
| Discoverability | 8 | 84% (+26%) | 68% (+15%) |
| Effectiveness | 8 | 90% (+5%) | 86% (-0%) |
| Efficiency | 8 | 61% (+24%) | 50% (+10%) |
Skill Version(s): <br>
92960d7 (source: git SHA, committed 2026-05-29) <br>
Ethical Considerations: <br>
NVIDIA believes Trustworthy AI is a shared responsibility and we have established policies and practices to enable development for a wide array of AI applications. When downloaded or used in accordance with our terms of service, developers should work with their internal team to ensure this skill meets requirements for the relevant industry and use case and addresses unforeseen product misuse. <br>
(For Release on NVIDIA Platforms Only) <br> Please report quality, risk, security vulnerabilities or NVIDIA AI Concerns here. <br>
Related skills
FAQ
What does accelerated-computing-cudf do?
accelerated-computing-cudf is a Claude Code skill for ai & agent building.
When should I use accelerated-computing-cudf?
When you need to helps with ai & agent building tasks., or when accelerated-computing-cudf is a claude code skill for ai & agent building.
What are the main capabilities?
accelerated-computing-cudf; AI & Agent Building; AI-coding skill.