Turn readings into events
Raw datapoints answer "what is the value?"; events answer "what happened?". A common pipeline reads recent readings, checks them against a rule, and records a discrete, queryable event when the rule fires — a threshold breach, a state change, an alarm. Events carry a type, a time, and metadata, and can reference the resources they concern.
Detect a threshold breach and record it
Read the latest hour, and if any reading exceeds a limit, create an event.
- Java
- Python
- Rust
import java.time.ZonedDateTime;
var filter = new RetrieveFilter();
filter.setExternalId("engine_temperature");
filter.setStart(ZonedDateTime.now().minusHours(1));
filter.setEnd(ZonedDateTime.now());
var request = new DataRetriever<RetrieveFilter>();
request.setItems(List.of(filter));
var series = client.timeseries().retrieve(request).getItems().get(0);
boolean tooHot = series.getDatapoints().stream()
.anyMatch(p -> Double.parseDouble(p.getValue()) > 110.0);
if (tooHot) {
EventModel event = new EventModel();
event.setExternalId("overheat_press_07_" + System.currentTimeMillis());
event.setType("overheat");
event.setStatus("open");
event.setMetadata(Map.of("series", "engine_temperature", "limit", "110"));
event.setEventTime(ZonedDateTime.now());
client.events().create(List.of(event));
}
import datahub_sdk, pandas as pd
rf = datahub_sdk.RetrieveFilter(
ts="engine_temperature",
start=pd.Timestamp.now(tz="UTC") - pd.Timedelta(hours=1),
end=pd.Timestamp.now(tz="UTC"))
points = client.timeseries.retrieve_datapoints(rf)[0].get_datapoints()
if any(float(dp.value) > 110.0 for dp in points):
client.events.create([datahub_sdk.Event(
external_id=f"overheat_press_07_{int(pd.Timestamp.now().timestamp())}",
type="overheat",
status="open",
event_time=pd.Timestamp.now(tz="UTC"),
metadata={"series": "engine_temperature", "limit": "110"})])
use chrono::Utc;
use dataplatform_rust_sdk::generic::{DataWrapper, RetrieveFilter};
use dataplatform_rust_sdk::events::Event;
let filter = RetrieveFilter {
external_id: Some("engine_temperature".into()),
start: Some(Utc::now() - chrono::Duration::hours(1)),
end: Some(Utc::now()),
..Default::default()
};
let series = api.time_series
.retrieve_datapoints(&DataWrapper::from(vec![filter])).await?
.get_items().remove(0);
let too_hot = series.datapoints.iter()
.any(|p| p.value.as_deref().and_then(|v| v.parse::<f64>().ok()).unwrap_or(0.0) > 110.0);
if too_hot {
let mut event = Event::new(format!("overheat_press_07_{}", Utc::now().timestamp()));
event.r#type = Some("overheat".into());
event.status = Some("open".into());
event.add_metadata("series".into(), "engine_temperature".into());
event.add_metadata("limit".into(), "110".into());
event.set_event_time(Utc::now());
api.events.create(&vec![event]).await?;
}
Query the events later
Events are first-class records — filter them by type, time or metadata for an audit trail or an incident timeline.
- Java
- Python
- Rust
EventRetreiver retriever = new EventRetreiver();
retriever.setLimit(100);
retriever.getFilter().setType("overheat");
DataWrapper<EventModel> overheats = client.events().filter(retriever);
overheats = client.events.filter(datahub_sdk.EventFilter(
basic_filter=datahub_sdk.BasicEventFilter(type="overheat"),
limit=100))
use dataplatform_rust_sdk::filters::{BasicEventFilter, EventFilter};
let filter = EventFilter::default()
.set_filter(BasicEventFilter { r#type: Some("overheat".into()), ..Default::default() })
.set_limit(100)
.build();
let overheats = api.events.filter(&filter).await?;
Run it on the live stream
Polling the last hour is fine for a cron job. To react the moment a reading crosses the line, run this rule inside a live subscription loop instead of on a timer.