Temporal events: the consumer contract¶
The temporal pipeline (pipeline.type: temporal, issues #12/#213) aggregates
events — spatiotemporal objects like storms — into one tabular row of scalar
attributes each. This page is the stable interface for libraries that feed
events into zagg.agg(): the event shapes both backends accept, the config
blocks that describe collections and specs, and the capability-resolution
policy that governs what runs where.
The event model¶
An event is one binary mask with dims (time, lat, lon): where and when
the object exists. process_event streams the mask's timesteps, applies a
per-timestep spatial_func under a mask provider, and reduces across time
with a temporal_reducer — one scalar per configured attribute.
There is deliberately no second mask channel: an attribute set needing a
different footprint (e.g. a precipitation lookahead window) is a different
event set. Run agg() once per mask family and join the tabular outputs on
event_key.
Passing events¶
zagg.agg(config, events=..., backend=...) accepts events in two shapes,
chosen by the backend:
Local backend — in-memory tuples, one per event:
(event_key, event_mask, collections, static_data)
event_key: hashable identifier, echoed into the output row.event_mask:xr.DataArray, dims(time, lat, lon).collections:{name: xr.Dataset}— the source fields, already opened. The runner applies the config's per-collection options (below) to each event's collections, so local and Lambda semantics match.static_data:{name: DataArray | Dataset}— e.g.ais_mask,cell_areas,climatology.
Lambda backend — URI dicts, one per event; the worker loads its own inputs, keeping the async payload far under Lambda's 256 KB Event cap:
{
"event_key": "storm_00042",
"event_mask_uri": "s3://.../masks/storm_00042.nc",
"collection_uris": {"merra2_slv": ["s3://.../day1.nc4", "s3://.../day2.nc4"]},
"static_uris": {"ais_mask": "s3://.../ais_mask.nc", "cell_areas": "s3://.../areas.nc"},
"s3_credentials": {...}, # optional: read creds for the collections
"input_credentials": "unsigned", # optional: mask/statics channel
}
event_mask_uri/static_urisvalues: a single URI each (local path,s3://;.zarropens as a store, a local non-.zarrpath opens directly, and ans3://non-.zarrobject is fetched as bytes and opened in-memory). A single-variable file becomes aDataArray.collection_urisvalues: one URI or a list — a multi-granule event (e.g. a storm spanning several MERRA-2 daily files) is opened per-URI, concatenated alongtime, and time-sorted.s3_credentials: optional; covers only the source collections it was fetched for. Events without it receive shared credentials fetched once from thedata_source.credentials_providerregistry name when the config sets one (nsidc,gesdisc, andlpdaacship as built-ins).input_credentials: optional; the channel for the consumer-owned mask and statics — an explicit creds dict,"unsigned"(anonymous requests: the correct mechanism for public buckets, since a request signed with scoped source credentials is denied cross-account even where anonymous access is granted), or absent for the worker's execution role.
The full worker payload (mode, config, return_results, result_url) is
assembled by the runner — consumers supply only the dicts above.
Describing collections¶
data_source.collections accepts a list of names, or a mapping carrying
declarative per-collection reader options:
data_source:
reader: xarray_s3
collections:
merra2_slv: # no options needed
merra2_precip:
coord_round: 5 # round lat/lon coords first (float dirt)
variables: [PRECCU, PRECLS, PRECSN] # subset before the rest
time_offset: "-30min" # shift stamps onto the hour
resample: {freq: "3h", how: sum, scale: 3600} # rates -> totals
derived:
rainfall: "PRECCU + PRECLS" # numpy expression
doi: "10.5067/Q5GVUVUIVGO7" # unknown keys pass through
Options apply in the order listed above (prepare_collection): coord_round
runs first (it rounds the source grid's own lat/lon coordinate arrays, so a
grid that ships float dirt still matches the rounded event/static coords),
then variables, time_offset, resample, and derived. derived
expressions evaluate in the same restricted namespace as the spatial
pipeline's expression fields: numpy plus the collection's variables, no
builtins. Unknown keys (like doi) are ignored by the reader — they are
metadata for catalog tooling. Validated option values fail at config load,
not per-worker.
Spec keys¶
Each aggregation.variables entry:
| key | required | meaning |
|---|---|---|
variable |
yes | variable name in the collection (may be derived) |
collection |
yes | collection name |
spatial_func |
yes | per-timestep reduction (registry name) |
temporal_reducer |
yes | cross-timestep accumulator (registry name) |
mask |
no (ais) |
mask provider: full / ais / ocean / registered |
anomaly |
no | sugar for transform: monthly_anomaly |
transform |
no | field-transform name applied per timestep |
negate |
no | negate the field (southward-positive fluxes) |
trigger |
no | event-trigger name gating which timesteps update (e.g. first_landfall) |
trigger_mask |
no | static field the trigger tests against (default ais_mask) |
params |
no | free-form mapping threaded into mask providers, field transforms, and event triggers (not spatial_funcs or reducers) |
Capability resolution and the worker policy¶
Every name above resolves through zagg.registry at run time. The Lambda
payload is pure data — capability names, numpy expressions, URIs — and a
name only resolves against what is installed in the worker's layer. There is
no mechanism for shipping third-party code to the fleet, by design (issue
213): on backend="lambda", worker-side capabilities are zagg built-ins and¶
declarative config, full stop.
The registries stay open on backend="local": register a custom trigger,
reducer, mask, transform, or reader (directly via zagg.registry.register_*
or a zagg.plugins entry point) and run it on your own machine. To make a
capability Lambda-eligible, upstream it into zagg as a built-in — prototyping
locally as a plugin and promoting what proves out is the intended
contribution funnel.
Credential providers run orchestrator-side, so they carry no such
restriction; data_source.credentials_provider may name a plugin provider on
either backend — including non-NASA S3-compatible sources. The spatial
pipeline honors the same key for its S3-driver source reads, defaulting to
nsidc when the key is absent; the HTTPS driver instead authenticates with an
EDL bearer token and does not consult credential providers.