radrs.raystack¶
Raystack format tools for ML training.
This module provides tools for converting NEXRAD data into the raystack
format — a flat layout of vcps, sweeps, returns, and activity arrays
optimized for machine learning training pipelines.
Returns are chunked segments of physical radials. A single radial can produce
multiple returns when n_gates > fold_size. Activity metrics are computed
over these return chunks.
Examples:
Open a NEXRAD volume as a raystack DataTree:
import radrs.raystack as rrs
rdt = rrs.open_datatree("s3://unidata-nexrad-level2/2024/03/15/KTLX/KTLX20240315_000217_V06", fold_size=128)
rdt["returns"]["DBZH"].shape # (n_returns, 128)
Convert an existing xradar DataTree to raystack:
import radrs.xradar as rxr
import radrs.raystack as rrs
dt = rxr.open_datatree(file_path)
rs = rrs.from_xradar_datatree(dt, fold_size=128)
Apply built-in QC at parse time:
import radrs.raystack as rrs
import radrs.qc as qc
rdt = rrs.open_datatree(src, fold_size=128, qc=[qc.RhohvThreshold(threshold=0.8)])
For lower-level use over already-loaded bytes, see parse.
BatchedRaystack ¶
Pre-allocated raystack batch accumulator for incremental filling.
Pre-allocates flat arrays for P patterns (VCPs), S sweeps, and R returns, then fills incrementally from volumes added to the batch. Each return tracks both its parent sweep_time and vcp_time for temporal organization.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
max_vcps
|
int
|
Maximum number of VCP patterns |
required |
max_sweeps
|
int
|
Maximum total number of sweeps across all patterns |
required |
max_returns
|
int
|
Maximum total number of returns/radials |
required |
fold_size
|
int
|
Range fold size |
128
|
truncate
|
bool
|
If True, output arrays are truncated to actual filled size. If False, output arrays remain at max_returns size with NaN/0 fill. |
True
|
drop_empty_returns
|
bool
|
If True, returns with all NaN moment data are excluded from the output. If False, all returns are included regardless of data completeness. |
False
|
include_sweeps
|
bool
|
If True, sweep-level metadata is included in the output. If False, only VCP metadata is included. |
True
|
include_returns
|
bool
|
If True, return-level data (radials and moment data) are included in the output. If False, only VCP and sweep metadata are collected. |
True
|
Examples:
Accumulate volumes from an L2 archive iterator into a batch:
>>> import radrs
>>> from datetime import datetime, timezone
>>> # Allocate for 10 VCPs, ~140 sweeps (14 per VCP), ~50k returns
>>> batch = radrs.raystack.BatchedRaystack(
... max_vcps=10,
... max_sweeps=140,
... max_returns=50000,
... fold_size=128,
... )
>>> archive = radrs.NexradL2ArchiveIter(
... base_uri="s3://unidata-nexrad-level2",
... start_time=datetime(2024, 3, 15, tzinfo=timezone.utc),
... end_time=datetime(2024, 3, 16, tzinfo=timezone.utc),
... storage_options={"anon": "true"},
... site_filter=["KTLX"],
... )
>>> n_added = batch.add_volumes_from_l2(archive, prefetch=8)
>>> prog = batch.progress()
>>> print(f"Collected {prog['patterns_filled']} patterns, "
... f"{prog['sweeps_filled']} sweeps, {prog['returns_filled']} returns")
Convert to DataTree for xarray:
Pre-allocate fixed-size arrays (no truncation):
>>> # Request 75000 returns, get exactly 75000-element arrays
>>> batch = radrs.raystack.BatchedRaystack(
... max_vcps=10,
... max_sweeps=140,
... max_returns=75000,
... fold_size=128,
... truncate=False, # Keep full size with NaN fill
... )
>>> n_added = batch.add_volumes_from_l2(archive, prefetch=8)
>>> raystack = batch.finalize_to_dict()
>>> print(raystack['returns']['azimuth'].shape) # (75000,) regardless of actual fills
add_volume ¶
Add a complete volume from raw bytes.
Processes the volume directly from raw NEXRAD file bytes, avoiding intermediate dict allocation for efficiency.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
data
|
bytes
|
Raw NEXRAD volume file bytes |
required |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
If capacity is exceeded or batch is finalized |
Examples:
add_volume_from_url ¶
Add a volume from a cloud or local URL.
Fetches the volume from the specified URL and adds it to the batch. Supports S3, GCS, Azure Blob Storage, and local filesystem URLs.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
url
|
str
|
Full URL to the volume file. Supported schemes: - S3: s3://bucket/path/to/file - GCS: gs://bucket/path/to/file - Azure: az://container/path/to/file or azure://container/path/to/file - Local: /path/to/file or file:///path/to/file |
required |
storage_options
|
dict
|
Storage backend configuration. Examples: - S3: {"region": "us-east-1", "anon": "true"} - GCS: {"service_account_path": "/path/to/key.json"} - Azure: {"account_name": "...", "access_key": "..."} |
None
|
Raises:
| Type | Description |
|---|---|
RuntimeError
|
If capacity is exceeded, batch is finalized, or fetch fails |
Examples:
S3 with anonymous access:
>>> stream = radrs.raystack.BatchedRaystack(10, 140, 50000)
>>> stream.add_volume_from_url(
... "s3://unidata-nexrad-level2/2024/03/15/KTLX/KTLX20240315_000217_V06",
... storage_options={"anon": "true"}
... )
GCS with service account:
>>> stream.add_volume_from_url(
... "gs://my-bucket/nexrad/data/KTLX20240315_120000_V06",
... storage_options={"service_account_path": "/path/to/key.json"}
... )
Local filesystem:
Azure Blob Storage:
add_volumes_from_l2 ¶
Add volumes from a NexradL2ArchiveIter with prefetch support.
This method supports multi-cloud sources (S3, GCS, Azure, local filesystem) using the NexradL2ArchiveIter which understands NEXRAD L2 Archive directory structure (YYYY/MM/DD/SITE/). Only volumes with timestamps within the specified time range are returned.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
l2_iter
|
NexradL2ArchiveIter
|
L2 archive iterator configured with time bounds |
required |
prefetch
|
int
|
Number of volumes to prefetch in parallel |
2
|
Returns:
| Type | Description |
|---|---|
int
|
Number of volumes successfully added |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
If batch is already finalized |
Notes
Running out of capacity is not an error. A volume that does not fit is rejected whole (adds are atomic, so a partial volume never lands) and iteration continues, so a later, smaller volume can still be admitted: an undersized batch quietly yields a time range with holes in it rather than a clean prefix. Iteration stops early only once a dimension is exactly full. Compare the returned count against the number of volumes you expected.
Fetch failures are skipped the same way, as are parse failures on
the default include_sweeps=True path; set RADRS_LOG=warn to
see the reason for each skip. With include_sweeps=False, a
volume that cannot be peeked raises instead of being skipped.
Examples:
S3 with anonymous access (time range within a single day):
>>> import radrs
>>> import radrs.raystack as rrs
>>> from datetime import datetime, timezone
>>> l2_iter = radrs.NexradL2ArchiveIter(
... base_uri="s3://noaa-nexrad-level2",
... start_time=datetime(2024, 3, 15, 10, 0, 0, tzinfo=timezone.utc),
... end_time=datetime(2024, 3, 15, 14, 0, 0, tzinfo=timezone.utc),
... storage_options={"anon": "true"},
... site_filter=["KTLX"]
... )
>>> stream = rrs.BatchedRaystack(10, 140, 50000)
>>> n_added = stream.add_volumes_from_l2(l2_iter, prefetch=8)
Time range spanning multiple days:
>>> l2_iter = radrs.NexradL2ArchiveIter(
... base_uri="s3://noaa-nexrad-level2",
... start_time=datetime(2024, 3, 15, 20, 0, 0, tzinfo=timezone.utc),
... end_time=datetime(2024, 3, 16, 4, 0, 0, tzinfo=timezone.utc),
... storage_options={"anon": "true"},
... site_filter=["KTLX"]
... )
GCS with service account:
>>> l2_iter = radrs.NexradL2ArchiveIter(
... base_uri="gs://my-bucket/nexrad",
... start_time=datetime(2024, 3, 15, 0, 0, 0, tzinfo=timezone.utc),
... end_time=datetime(2024, 3, 16, tzinfo=timezone.utc),
... storage_options={"service_account_path": "/path/to/key.json"}
... )
>>> n_added = stream.add_volumes_from_l2(l2_iter, prefetch=8)
Local filesystem:
progress ¶
Get current fill progress.
Call this before finalize_to_dict() or finalize_to_rs_dt(),
which consume the batch — progress() panics once the buffers
have been handed to Python. Plain finalize() leaves these
figures untouched.
Notes
returns_filled counts the returns that survived
drop_empty_returns, while returns_capacity is checked against
the uncompacted count. With compaction on, fill_fraction is
therefore a lower bound on real capacity pressure — use the volume
count returned by add_volumes_from_l2 to confirm everything fit.
Returns:
| Type | Description |
|---|---|
dict
|
Progress with keys: patterns_filled, patterns_capacity, sweeps_filled, sweeps_capacity, returns_filled, returns_capacity, fill_fraction |
has_capacity ¶
Check if batch has remaining capacity.
Returns:
| Type | Description |
|---|---|
bool
|
True if batch can accept more data |
add_qc_outputs ¶
Add quality control outputs to the batch.
QC operations are computed over the data that has been accumulated so far. Call this before finalize() to include QC outputs in the final result.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
qc_steps
|
list of QCStep or QCStep
|
Quality control steps to apply. Examples: - radrs.qc.RhohvThreshold(threshold=0.8) - radrs.qc.SunSpike(dbzh_threshold=50.0) - radrs.qc.VradhWindingNumber(nyquist=25.0) |
required |
Examples:
>>> import radrs.raystack as rrs
>>> import radrs.qc as qc
>>> batch = rrs.BatchedRaystack(10, 140, 50000)
>>> # ... add volumes ...
>>> batch.add_qc_outputs([
... qc.RhohvThreshold(threshold=0.8, vname="rhohv_mask"),
... qc.SunSpike(vname="sun_spike")
... ])
>>> result = batch.finalize_to_dict()
>>> # QC outputs will be in result['returns'] with 'qc.' prefix
>>> print(result['returns']['qc.rhohv_mask'].shape)
finalize ¶
Finalize the batch in place.
With the default truncate=True the arrays are already at their
filled size and this is a no-op; with truncate=False they are
padded out to full capacity with NaN/NaT fill. Either way the
progress() counts are unchanged.
Called automatically by finalize_to_dict() and finalize_to_rs_dt().
finalize_to_dict ¶
Returns the batch as a raystack dict with all data zero-copied into python arrays.
Returns:
| Type | Description |
|---|---|
dict
|
Raystack dict with keys: vcps, sweeps, returns. Returns dict includes: vcp_time, sweep_time (parent times) |
finalize_to_rs_dt ¶
Convert to xarray Raystack DataTree.
Returns:
| Type | Description |
|---|---|
DataTree
|
Raystack DataTree with hierarchical structure |
open_datatree ¶
open_datatree(
source,
fold_size=None,
qc=None,
include_activity=True,
storage_options=None,
*,
format: str = DEFAULT_SOURCE_FORMAT,
)
Open a NEXRAD volume and return a raystack DataTree.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
source
|
str or bytes
|
Path to file, URI ( |
required |
fold_size
|
int
|
Range fold size (chunks each radial into segments). |
None
|
qc
|
list of QCStep or QCStep
|
Quality control steps to apply during parse. |
None
|
include_activity
|
bool
|
Include activity metrics in the output. |
True
|
storage_options
|
dict
|
Storage backend configuration forwarded to object_store. Examples: - S3 anonymous: {"anon": "true"} - GCS service account: {"service_account_path": "/path/to/key.json"} - Azure: {"account_name": "...", "access_key": "..."} |
None
|
open_datatree_async
async
¶
open_datatree_async(
source,
fold_size=None,
qc=None,
include_activity=True,
storage_options=None,
*,
format: str = DEFAULT_SOURCE_FORMAT,
)
Async variant of open_datatree.
See open_datatree for parameter descriptions, including the
single-volume-only caveat for non-S3 cloud sources.