From c9690d144112295c1a8a88347374a3bab4975909 Mon Sep 17 00:00:00 2001 From: changetheway Date: Fri, 24 Jul 2026 02:08:46 +0000 Subject: [PATCH 1/2] feat(profiler): support Ascend NPU profiling via torch_npu.profiler --- src/cache_dit/profiler.py | 97 +++++++++++++++++++++++++++++------- tests/utils/test_profiler.py | 44 ++++++++++++++++ 2 files changed, 123 insertions(+), 18 deletions(-) create mode 100644 tests/utils/test_profiler.py diff --git a/src/cache_dit/profiler.py b/src/cache_dit/profiler.py index 80518b09..4e660b55 100644 --- a/src/cache_dit/profiler.py +++ b/src/cache_dit/profiler.py @@ -3,14 +3,14 @@ Reference: Adapted from https://github.com/sgl-project/sglang/blob/main/python/sglang/bench_one_batch.py """ +import gzip import logging import os import time from pathlib import Path -from typing import List, Optional +from typing import List, Optional, Tuple import torch -from torch.profiler import ProfilerActivity, profile from .platforms import current_platform logger = logging.getLogger(__name__) @@ -19,11 +19,46 @@ PROFILER_DIR = os.getenv("CACHE_DIT_TORCH_PROFILER_DIR", "/tmp/cache_dit_profiles") +def _supported_device() -> bool: + """Whether the profiler supports the current device. + + CUDA devices are profiled with ``torch.profiler`` and Ascend NPU devices with + ``torch_npu.profiler``; other device types are not supported. + + :returns: True if the current device is a supported accelerator. + """ + + if not current_platform.is_accelerator_available(): + return False + return current_platform.device_type in ("cuda", "npu") + + +def _resolve_profiler_backend() -> Tuple: + """Select the profiler backend for the current platform. + + ``torch_npu.profiler`` mirrors the ``torch.profiler`` API, so both backends are + driven the same way; only the ``ProfilerActivity`` enum and the accelerator + activity name differ. + + :returns: A tuple of (profile callable, ProfilerActivity enum, accelerator activity name). + """ + + if current_platform.device_type == "npu": + import torch_npu.profiler as npu_profiler # noqa: WPS433 + + return npu_profiler.profile, npu_profiler.ProfilerActivity, "NPU" + from torch.profiler import ProfilerActivity, profile # noqa: WPS433 + + return profile, ProfilerActivity, "CUDA" + + class ProfilerContext: - """Context manager wrapper around `torch.profiler` for cache-dit runs. + """Context manager wrapper around ``torch.profiler`` / ``torch_npu.profiler``. - It centralizes trace-file naming, optional CUDA memory-history capture, and multi-rank output - layout so profiling can be enabled consistently from scripts or helper decorators. + It centralizes trace-file naming, optional accelerator memory-history capture, and + multi-rank output layout so profiling can be enabled consistently from scripts or + helper decorators. CUDA devices capture memory snapshots via ``torch.cuda.memory``; + Ascend NPU devices record memory through the profiler's own ``profile_memory`` option. """ def __init__( @@ -38,21 +73,27 @@ def __init__( """Configure a profiler session. :param enabled: Whether profiling should actually be activated. - :param activities: Activity names such as `CPU`, `GPU`, or `MEM`. + :param activities: Activity names such as `CPU`, `GPU`, or `MEM`. `GPU` maps to the + platform accelerator activity (CUDA on CUDA devices, NPU on Ascend devices). :param output_dir: Directory where traces and memory snapshots are written. :param profile_name: Base name used for profiler output files. :param with_stack: Whether to capture Python stacks for profiled ops. :param record_shapes: Whether to record tensor shapes in the profiler trace. """ - assert (current_platform.is_accelerator_available() and current_platform.device_type - == "cuda"), "Torch ProfilerContext currently only supports CUDA devices." + assert _supported_device(), ( + "Torch ProfilerContext currently only supports CUDA or Ascend NPU devices, " + f"got device_type={current_platform.device_type!r}.") self.enabled = enabled self.activities = activities or ["CPU", "GPU"] self.output_dir = Path(output_dir or PROFILER_DIR).expanduser() self.profile_name = profile_name or f"profile_{int(time.time())}" self.with_stack = with_stack self.record_shapes = record_shapes + # NPU records memory through the profiler (profile_memory=True); CUDA uses the + # separate torch.cuda.memory history mechanism instead. + self._is_npu = current_platform.device_type == "npu" + self._track_memory = "MEM" in self.activities self.profiler = None self.trace_path = None @@ -62,14 +103,20 @@ def __enter__(self): if not self.enabled: return self - assert (current_platform.is_accelerator_available() and current_platform.device_type - == "cuda"), "Torch ProfilerContext currently only supports CUDA devices." + assert _supported_device(), ( + "Torch ProfilerContext currently only supports CUDA or Ascend NPU devices, " + f"got device_type={current_platform.device_type!r}.") self.output_dir.mkdir(parents=True, exist_ok=True) + profile_fn, activity_enum, accelerator_activity = _resolve_profiler_backend() + # "GPU" maps to the platform accelerator activity (CUDA on CUDA, NPU on NPU); + # the accelerator name (e.g. "NPU") is also accepted as an explicit synonym. + accelerator_activity_enum = getattr(activity_enum, accelerator_activity) activity_map = { - "CPU": ProfilerActivity.CPU, - "GPU": ProfilerActivity.CUDA, + "CPU": activity_enum.CPU, + "GPU": accelerator_activity_enum, + accelerator_activity: accelerator_activity_enum, } torch_activities = [activity_map[a] for a in self.activities if a in activity_map] @@ -85,16 +132,19 @@ def __enter__(self): filename = "-".join(filename_parts) + ".trace.json.gz" self.trace_path = self.output_dir / filename - if "MEM" in self.activities and torch.cuda.is_available(): + if self._track_memory and not self._is_npu and torch.cuda.is_available(): torch.cuda.memory._record_memory_history(max_entries=100000) logger.info("Started CUDA memory profiling") if torch_activities: - self.profiler = profile( + profiler_kwargs = dict( activities=torch_activities, with_stack=self.with_stack, record_shapes=self.record_shapes, ) + if self._is_npu and self._track_memory: + profiler_kwargs["profile_memory"] = True + self.profiler = profile_fn(**profiler_kwargs) self.profiler.start() logger.info(f"Started profiling. Traces will be saved to: {self.output_dir} " @@ -107,17 +157,28 @@ def __exit__(self, exc_type, exc_val, exc_tb): return if self.profiler is not None: - if torch.cuda.is_available(): - torch.cuda.synchronize() + if current_platform.is_accelerator_available(): + current_platform.synchronize() self.profiler.stop() logger.info(f"Exporting trace to: {self.trace_path}") - self.profiler.export_chrome_trace(str(self.trace_path)) + if self._is_npu: + # torch_npu.profiler.export_chrome_trace can only write a plain .json file; + # export to a temporary json then gzip it so the .trace.json.gz artifact + # convention used on CUDA is preserved. with_suffix("") strips the trailing + # ".gz" so the temp path keeps a ".json" suffix the exporter accepts. + tmp_trace = self.trace_path.with_suffix("") + self.profiler.export_chrome_trace(str(tmp_trace)) + with open(tmp_trace, "rb") as src, gzip.open(self.trace_path, "wb") as dst: + dst.writelines(src) + tmp_trace.unlink(missing_ok=True) + else: + self.profiler.export_chrome_trace(str(self.trace_path)) logger.info(f"Profiling completed. Trace saved to: {self.trace_path}") - if "MEM" in self.activities and torch.cuda.is_available(): + if self._track_memory and not self._is_npu and torch.cuda.is_available(): timestamp = int(time.time()) rank = torch.distributed.get_rank() if torch.distributed.is_initialized() else 0 memory_snapshot_path = (self.output_dir / diff --git a/tests/utils/test_profiler.py b/tests/utils/test_profiler.py new file mode 100644 index 00000000..b1000267 --- /dev/null +++ b/tests/utils/test_profiler.py @@ -0,0 +1,44 @@ +import pytest + +from cache_dit.platforms import current_platform +from cache_dit.profiler import ( + ProfilerContext, + _resolve_profiler_backend, + _supported_device, +) + + +class TestProfilerBackendSelection: + + def test_supported_device_returns_bool(self): + assert isinstance(_supported_device(), bool) + + def test_supported_device_matches_platform(self): + # CPU is not a supported profiler device; CUDA/NPU are. + if current_platform.device_type == "cpu": + assert _supported_device() is False + else: + assert _supported_device() is current_platform.is_accelerator_available() + + def test_resolve_backend_returns_cuda_on_non_npu(self): + # On non-NPU devices the CUDA profiler backend is selected; NPU devices select + # torch_npu.profiler. Either way the tuple shape is (callable, enum, str). + profile_fn, activity_enum, accelerator = _resolve_profiler_backend() + assert callable(profile_fn) + assert hasattr(activity_enum, "CPU") + assert isinstance(accelerator, str) + if current_platform.device_type == "npu": + assert accelerator == "NPU" + else: + assert accelerator == "CUDA" + + +class TestProfilerContextGuard: + + def test_construct_raises_on_unsupported_device(self): + # ProfilerContext must refuse unsupported devices (e.g. CPU) instead of + # silently producing an empty/invalid trace. + if _supported_device(): + pytest.skip("device is supported by the profiler") + with pytest.raises(AssertionError): + ProfilerContext(enabled=True) From f4084f956684a397896b429230464dd9aa385a41 Mon Sep 17 00:00:00 2001 From: changetheway Date: Fri, 24 Jul 2026 06:39:43 +0000 Subject: [PATCH 2/2] feat(profiler): export full CANN output directory on NPU Switch the NPU path from export_chrome_trace (single json) to on_trace_ready=tensorboard_trace_handler + experimental_config, so the full Ascend profiling output directory is produced: ASCEND_PROFILER_OUTPUT (kernel_details.csv, step_trace_time.csv, trace_view.json, op_summary.csv, op_statistic.csv, operator_details.csv, .db) plus PROF_XXX, FRAMEWORK and logs. data_simplification=False keeps every generated file. - profiler_level configurable via CACHE_DIT_NPU_PROFILER_LEVEL (Level0 default) - sanitize worker_name to the allowed [A-Za-z0-9_-] charset - locate the produced *_ascend_pt dir and expose it via trace_path - CUDA path unchanged --- src/cache_dit/profiler.py | 194 +++++++++++++++++++++++++++++------ tests/utils/test_profiler.py | 19 ++++ 2 files changed, 182 insertions(+), 31 deletions(-) diff --git a/src/cache_dit/profiler.py b/src/cache_dit/profiler.py index 4e660b55..3a2211d7 100644 --- a/src/cache_dit/profiler.py +++ b/src/cache_dit/profiler.py @@ -3,9 +3,10 @@ Reference: Adapted from https://github.com/sgl-project/sglang/blob/main/python/sglang/bench_one_batch.py """ -import gzip +import glob import logging import os +import re import time from pathlib import Path from typing import List, Optional, Tuple @@ -18,6 +19,11 @@ # Default profiler directory PROFILER_DIR = os.getenv("CACHE_DIT_TORCH_PROFILER_DIR", "/tmp/cache_dit_profiles") +# NPU profiling depth (Level0 | Level1 | Level2 | none). Higher levels collect more +# CANN/AscendCL data (e.g. communication.json, api_statistic.csv, aic metrics) at the +# cost of larger output and slower parsing. See the Ascend PyTorch Profiler docs. +NPU_PROFILER_LEVEL_ENV = "CACHE_DIT_NPU_PROFILER_LEVEL" + def _supported_device() -> bool: """Whether the profiler supports the current device. @@ -52,13 +58,92 @@ def _resolve_profiler_backend() -> Tuple: return profile, ProfilerActivity, "CUDA" +def _npu_profiler_module(): + """Import and return the ``torch_npu.profiler`` module (NPU only).""" + + import torch_npu.profiler as npu_profiler # noqa: WPS433 + + return npu_profiler + + +def _npu_profiler_level(npu_profiler): + """Resolve the NPU profiler level from the env (default Level0). + + :param npu_profiler: The ``torch_npu.profiler`` module. + :returns: A ``ProfilerLevel`` enum value. + """ + + raw = os.getenv(NPU_PROFILER_LEVEL_ENV, "Level0").strip().lower() + table = { + "none": npu_profiler.ProfilerLevel.Level_none, + "level_none": npu_profiler.ProfilerLevel.Level_none, + "0": npu_profiler.ProfilerLevel.Level0, + "level0": npu_profiler.ProfilerLevel.Level0, + "1": npu_profiler.ProfilerLevel.Level1, + "level1": npu_profiler.ProfilerLevel.Level1, + "2": npu_profiler.ProfilerLevel.Level2, + "level2": npu_profiler.ProfilerLevel.Level2, + } + return table.get(raw, npu_profiler.ProfilerLevel.Level0) + + +def _npu_experimental_config(npu_profiler): + """Build the NPU experimental config that produces the full output directory. + + ``export_type=Text`` emits the ``.json``/``.csv`` timeline and summary files plus the + aggregate ``.db`` files; ``data_simplification=False`` keeps every generated file + (kernel_details.csv, step_trace_time.csv, trace_view.json, op_summary.csv, the + ``PROF_XXX`` raw data, ``FRAMEWORK``/``logs``) instead of pruning them. + + :param npu_profiler: The ``torch_npu.profiler`` module. + :returns: A tuple of (experimental config, resolved profiler level). + """ + + level = _npu_profiler_level(npu_profiler) + # aic_metrics must match the level, otherwise torch_npu resets it with a warning. + aic_metrics = (npu_profiler.AiCMetrics.PipeUtilization if level + in (npu_profiler.ProfilerLevel.Level1, + npu_profiler.ProfilerLevel.Level2) else npu_profiler.AiCMetrics.AiCoreNone) + config = npu_profiler._ExperimentalConfig( + export_type=[npu_profiler.ExportType.Text], + profiler_level=level, + aic_metrics=aic_metrics, + data_simplification=False, + ) + return config, level + + +def _npu_worker_name(profile_name: str, rank: int, world_size: int) -> str: + """Build a valid NPU ``worker_name`` from the profile name. + + ``worker_name`` only allows letters, digits, underscores and hyphens, and becomes the + ``{worker_name}_{timestamp}_ascend_pt`` output directory name. + + :param profile_name: The base profile name. + :param rank: The distributed rank. + :param world_size: The distributed world size. + :returns: A sanitized worker name. + """ + + safe = re.sub(r"[^A-Za-z0-9_-]", "_", str(profile_name or "cache_dit")).strip("_") + safe = safe or "cache_dit" + if world_size > 1: + safe = f"{safe}-rank{rank}" + return safe + + class ProfilerContext: """Context manager wrapper around ``torch.profiler`` / ``torch_npu.profiler``. It centralizes trace-file naming, optional accelerator memory-history capture, and multi-rank output layout so profiling can be enabled consistently from scripts or - helper decorators. CUDA devices capture memory snapshots via ``torch.cuda.memory``; - Ascend NPU devices record memory through the profiler's own ``profile_memory`` option. + helper decorators. + + On CUDA it exports a gzip chrome trace and optional ``torch.cuda.memory`` snapshots. + On Ascend NPU it drives ``on_trace_ready=tensorboard_trace_handler`` so the full CANN + profiling output directory is produced (``ASCEND_PROFILER_OUTPUT`` with + ``kernel_details.csv`` / ``step_trace_time.csv`` / ``trace_view.json`` / ``op_summary.csv`` + plus the ``PROF_XXX`` raw data, ``FRAMEWORK`` and ``logs``), not just a single trace. """ def __init__( @@ -94,6 +179,8 @@ def __init__( # separate torch.cuda.memory history mechanism instead. self._is_npu = current_platform.device_type == "npu" self._track_memory = "MEM" in self.activities + # NPU: worker name used for the {worker_name}_{ts}_ascend_pt output directory. + self._npu_worker = None self.profiler = None self.trace_path = None @@ -126,31 +213,62 @@ def __enter__(self): rank = torch.distributed.get_rank() world_size = torch.distributed.get_world_size() + if not torch_activities: + return self + + if self._is_npu: + self._start_npu_profiler(profile_fn, rank, world_size, torch_activities) + else: + self._start_cuda_profiler(profile_fn, rank, world_size, torch_activities) + + return self + + def _start_npu_profiler(self, profile_fn, rank, world_size, torch_activities): + """Configure and start the NPU profiler with full-directory export.""" + + npu_profiler = _npu_profiler_module() + self._npu_worker = _npu_worker_name(self.profile_name, rank, world_size) + experimental_config, level = _npu_experimental_config(npu_profiler) + trace_handler = npu_profiler.tensorboard_trace_handler( + dir_name=str(self.output_dir), + worker_name=self._npu_worker, + analyse_flag=True, + async_mode=False, + ) + self.profiler = profile_fn( + activities=torch_activities, + with_stack=self.with_stack, + record_shapes=self.record_shapes, + profile_memory=self._track_memory, + on_trace_ready=trace_handler, + experimental_config=experimental_config, + ) + self.profiler.start() + logger.info(f"Started NPU profiling. Full CANN output will be saved under " + f"{self.output_dir}/{self._npu_worker}__ascend_pt " + f"(profiler_level={level}, activities: {self.activities})") + + def _start_cuda_profiler(self, profile_fn, rank, world_size, torch_activities): + """Configure and start the CUDA profiler (gzip chrome trace).""" + filename_parts = [self.profile_name] if world_size > 1: filename_parts.append(f"rank{rank}") filename = "-".join(filename_parts) + ".trace.json.gz" self.trace_path = self.output_dir / filename - if self._track_memory and not self._is_npu and torch.cuda.is_available(): + if self._track_memory and torch.cuda.is_available(): torch.cuda.memory._record_memory_history(max_entries=100000) logger.info("Started CUDA memory profiling") - if torch_activities: - profiler_kwargs = dict( - activities=torch_activities, - with_stack=self.with_stack, - record_shapes=self.record_shapes, - ) - if self._is_npu and self._track_memory: - profiler_kwargs["profile_memory"] = True - self.profiler = profile_fn(**profiler_kwargs) - - self.profiler.start() - logger.info(f"Started profiling. Traces will be saved to: {self.output_dir} " - f"(activities: {self.activities})") - - return self + self.profiler = profile_fn( + activities=torch_activities, + with_stack=self.with_stack, + record_shapes=self.record_shapes, + ) + self.profiler.start() + logger.info(f"Started profiling. Traces will be saved to: {self.output_dir} " + f"(activities: {self.activities})") def __exit__(self, exc_type, exc_val, exc_tb): if not self.enabled: @@ -162,21 +280,13 @@ def __exit__(self, exc_type, exc_val, exc_tb): self.profiler.stop() - logger.info(f"Exporting trace to: {self.trace_path}") if self._is_npu: - # torch_npu.profiler.export_chrome_trace can only write a plain .json file; - # export to a temporary json then gzip it so the .trace.json.gz artifact - # convention used on CUDA is preserved. with_suffix("") strips the trailing - # ".gz" so the temp path keeps a ".json" suffix the exporter accepts. - tmp_trace = self.trace_path.with_suffix("") - self.profiler.export_chrome_trace(str(tmp_trace)) - with open(tmp_trace, "rb") as src, gzip.open(self.trace_path, "wb") as dst: - dst.writelines(src) - tmp_trace.unlink(missing_ok=True) + self.trace_path = self._locate_npu_output() + logger.info(f"Profiling completed. NPU profile saved to: {self.trace_path}") else: + logger.info(f"Exporting trace to: {self.trace_path}") self.profiler.export_chrome_trace(str(self.trace_path)) - - logger.info(f"Profiling completed. Trace saved to: {self.trace_path}") + logger.info(f"Profiling completed. Trace saved to: {self.trace_path}") if self._track_memory and not self._is_npu and torch.cuda.is_available(): timestamp = int(time.time()) @@ -193,6 +303,28 @@ def __exit__(self, exc_type, exc_val, exc_tb): f.write(torch.cuda.memory_summary()) logger.info(f"Memory summary saved to: {memory_summary_path}") + def _locate_npu_output(self) -> Path: + """Find the ``{worker_name}_{timestamp}_ascend_pt`` directory produced by the handler. + + Parsing is synchronous (async_mode=False), so it is normally ready right after + ``stop()``; a short poll covers slow disk finalization. + + :returns: The produced output directory (falls back to ``output_dir`` if not found). + """ + + pattern = str(self.output_dir / f"{self._npu_worker}_*_ascend_pt") + for _ in range(10): + ready = [ + c for c in glob.glob(pattern) if os.path.isdir(os.path.join(c, "ASCEND_PROFILER_OUTPUT")) + ] + if ready: + return Path(max(ready, key=os.path.getmtime)) + time.sleep(0.5) + candidates = glob.glob(pattern) + if candidates: + return Path(max(candidates, key=os.path.getmtime)) + return self.output_dir + def profile_function( enabled: bool = True, diff --git a/tests/utils/test_profiler.py b/tests/utils/test_profiler.py index b1000267..1a3fc559 100644 --- a/tests/utils/test_profiler.py +++ b/tests/utils/test_profiler.py @@ -3,6 +3,7 @@ from cache_dit.platforms import current_platform from cache_dit.profiler import ( ProfilerContext, + _npu_worker_name, _resolve_profiler_backend, _supported_device, ) @@ -33,6 +34,24 @@ def test_resolve_backend_returns_cuda_on_non_npu(self): assert accelerator == "CUDA" +class TestNpuWorkerName: + + def test_sanitizes_disallowed_chars(self): + # worker_name only allows [A-Za-z0-9_-]; dots and spaces must be replaced. + assert _npu_worker_name("wan2.2_t2v profile", 0, 1) == "wan2_2_t2v_profile" + + def test_appends_rank_for_distributed(self): + assert _npu_worker_name("flux", 3, 8) == "flux-rank3" + + def test_single_rank_omits_suffix(self): + assert _npu_worker_name("flux", 0, 1) == "flux" + + def test_falls_back_for_empty_name(self): + assert _npu_worker_name("", 0, 1) == "cache_dit" + # All-disallowed chars sanitize to empty and fall back to the default. + assert _npu_worker_name("...", 0, 1) == "cache_dit" + + class TestProfilerContextGuard: def test_construct_raises_on_unsupported_device(self):