Skip to content

Commit 7dfebba

Browse files
authored
Improved OCEL Import and Memory Usage + Misc Bug Fixes (#57)
* Initial version of OCEL Appendable/Readable trait for streaming-like memory-efficient import/export * More exposed binding functions + SlimLinkedOCEL Tweaks * Add qualifier-optional relation getters * Fix SQL-based export (timezone handling + float precision) * Additional r4pm bindings for KPIs/OC-Performance Analysis * Update bindings (non-aggregated performance analysis) * Fix double serialization for bindings (instead stick with just JSON Vec) * Update oc-perf functions to return top-k durations * Split SlimLinkedOCEL binding functions, fix SQL * Add oc_statistics module * Address review comments
1 parent 8f39ebe commit 7dfebba

38 files changed

Lines changed: 3854 additions & 1231 deletions

File tree

CHANGELOG.md

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,23 @@
11
# Changelog
22

3+
## Unreleased
4+
5+
- New `ReadableOCEL` and `AppendableOCEL` traits; OCEL exporters (CSV, XML, SQL, JSON) and the JSON importer are generic over them
6+
- `SlimLinkedOCEL` implements `AppendableOCEL`, with auto-declare types, auto-grow attributes, and value coercion via new `OCELAttributeValue::try_coerce_to`; misordered streams are buffered and resolved on `finalize`
7+
- `import_ocel_json_into` and `import_ocel_xml_into` stream JSON / XML directly into any `AppendableOCEL` (e.g., `SlimLinkedOCEL`) without materializing an `OCEL` first; `SlimLinkedOCEL::import_from_*` uses the streaming paths automatically
8+
- `ReadableOCEL::iter_events_of_type` / `iter_objects_of_type` (default filters; `SlimLinkedOCEL` overrides via per-type indices, avoiding a full scan per type in SQL export)
9+
- New `OCELAttributeType::as_type_str` returns `&'static str` (cheaper than `to_type_string`)
10+
- `EventIndex` / `ObjectIndex` are now `u32`-backed; `into_inner` returns `u32` (**Breaking**)
11+
- `SlimLinkedOCEL`, `SlimOCELEvent`, `SlimOCELObject` no longer derive `Deserialize` (still serialize); construct via `Importable` or `AppendableOCEL` (**Breaking**)
12+
- `SlimOCELEvent::relationships` and `SlimOCELObject::relationships` are now `Vec<(QualifierIdx, ObjectIndex)>` (was `Vec<(String, ObjectIndex)>`); resolve qualifier strings via `SlimLinkedOCEL::qualifier_str` (**Breaking**)
13+
- CSV exporter streams rows instead of buffering all rows; tracks `(time, object_id)` pairs in a `HashSet` to skip redundant object-attribute rows
14+
- `SlimLinkedOCEL::from_ocel` now goes through `AppendableOCEL`, so it shares attribute coercion, schema-grow, and duplicate-id detection with the streaming import paths. Behavior changes: events/objects with attribute names not in the declared type schema grow the schema (before they were silently dropped); duplicate event/object ids are skipped with a warning (before they were kept but unreachable); references to undeclared types auto-create the type (before they were dropped with warning)
15+
- `try_json_to_ocel` returning `Result<OCEL, serde_json::Error>` added alongside the existing panicking `json_to_ocel`
16+
- `SlimLinkedOCEL::get_o2o_rev_obs_of_obtype` / `get_e2o_rev_evs_of_evtype`: qualifier-optional reverse-relation getters filtered by type
17+
- New public object-centric analysis functions (also exposed as bindings): per-event sojourn and synchronization times with optional `top_k` (`analysis::object_centric::oc_performance`), E2O `(event_type, object_type)` counts and `source -> target` conversion rate (`analysis::object_centric::oc_statistics`), and per-object-type directly-follows graph and activity-trace variants (`discovery::object_centric::dfg` / `variants`)
18+
- Fix SQL export/import of floats and timestamps: floats are written as `DOUBLE PRECISION` (full f64 precision) and timestamps as naive UTC (avoids a double-applied timezone offset); import maps `DOUBLE` / `DOUBLE PRECISION` columns back to float, so round-trips no longer drop float attributes
19+
- New direct dependency on `hashbrown` for the slim per-id hash tables
20+
321
## 0.5.6
422

523
- Translate a `ProcessTree` into a `PetriNet`:

Cargo.lock

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

macros_process_mining/src/lib.rs

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -325,16 +325,16 @@ pub fn register_binding(args: TokenStream, item: TokenStream) -> TokenStream {
325325
let serialization_logic = if attrs.debug_output {
326326
quote! {
327327
let final_result = format!("{:?}", result);
328-
serde_json::to_value(final_result).map_err(|e| e.to_string())
328+
serde_json::to_vec(&final_result).map_err(|e| e.to_string())
329329
}
330330
} else if attrs.stringify_error {
331331
quote! {
332332
let ok_result = result.map_err(|e| e.to_string())?;
333-
serde_json::to_value(ok_result).map_err(|e| e.to_string())
333+
serde_json::to_vec(&ok_result).map_err(|e| e.to_string())
334334
}
335335
} else {
336336
quote! {
337-
serde_json::to_value(result).map_err(|e| e.to_string())
337+
serde_json::to_vec(&result).map_err(|e| e.to_string())
338338
}
339339
};
340340

@@ -398,7 +398,7 @@ pub fn register_binding(args: TokenStream, item: TokenStream) -> TokenStream {
398398
quote! {
399399
let id = format!("res_{}", uuid::Uuid::new_v4());
400400
__state_guard.insert(id.clone(), crate::bindings::RegistryItem::#variant_ident(result));
401-
serde_json::to_value(id).map_err(|e| e.to_string())
401+
serde_json::to_vec(&id).map_err(|e| e.to_string())
402402
}
403403
} else {
404404
serialization_logic.clone()
@@ -421,7 +421,7 @@ pub fn register_binding(args: TokenStream, item: TokenStream) -> TokenStream {
421421
};
422422
let id = format!("res_{}", uuid::Uuid::new_v4());
423423
state_lock.add(&id, crate::bindings::RegistryItem::#variant_ident(result));
424-
serde_json::to_value(id).map_err(|e| e.to_string())
424+
serde_json::to_vec(&id).map_err(|e| e.to_string())
425425
}
426426
} else {
427427
quote! {
@@ -468,7 +468,7 @@ pub fn register_binding(args: TokenStream, item: TokenStream) -> TokenStream {
468468
use serde_json::Value;
469469
use std::sync::RwLock;
470470

471-
fn #wrapper_name(args: &Value, state_lock: &AppState) -> Result<Value, String> {
471+
fn #wrapper_name(args: &Value, state_lock: &AppState) -> Result<Vec<u8>, String> {
472472
let arg_map = args.as_object().ok_or("Args must be JSON object")?;
473473
#execution_block
474474
}

process_mining/Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ chrono = { version = "0.4.40", features = ["serde"] }
1919
duckdb = { version = "1.2.1", optional = true, features = ["chrono"]}
2020
flate2 = "1.1.1"
2121
graphviz-rust = { version = "0.9.3", optional = true }
22+
hashbrown = "0.15"
2223
itertools = { version = "0.14.0" }
2324
kuzu = {version = "=0.11.2", optional = true}
2425
nalgebra = { version = "0.33.2", optional = true }
Lines changed: 28 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,11 @@
1-
use process_mining::{Importable, OCEL};
1+
use process_mining::core::event_data::object_centric::linked_ocel::{
2+
LinkedOCELAccess, SlimLinkedOCEL,
3+
};
4+
use process_mining::Importable;
25
use std::env;
36
use std::error::Error;
47
use std::path::PathBuf;
8+
use std::time::Instant;
59

610
fn main() -> Result<(), Box<dyn Error>> {
711
let args: Vec<String> = env::args().collect();
@@ -12,25 +16,32 @@ fn main() -> Result<(), Box<dyn Error>> {
1216

1317
let path = PathBuf::from(&args[1]);
1418
println!("Importing OCEL from {:?}", path);
19+
let now = Instant::now();
20+
let ocel = SlimLinkedOCEL::import_from_path(&path)?;
21+
println!("Successfully imported OCEL in {:?}.", now.elapsed());
22+
println!("Number of events: {}", ocel.get_all_evs().count());
23+
println!("Number of objects: {}", ocel.get_all_obs().count());
1524

16-
let ocel = OCEL::import_from_path(&path)?;
17-
println!("Successfully imported OCEL.");
18-
println!("Number of events: {}", ocel.events.len());
19-
println!("Number of objects: {}", ocel.objects.len());
20-
21-
println!(
22-
"Event Types: {:?}",
23-
ocel.event_types
24-
.iter()
25-
.map(|et| &et.name)
26-
.collect::<Vec<_>>()
27-
);
25+
println!("Event Types: {:?}", ocel.get_ev_types().collect::<Vec<_>>());
2826
println!(
2927
"Object Types: {:?}",
30-
ocel.object_types
31-
.iter()
32-
.map(|ot| &ot.name)
33-
.collect::<Vec<_>>()
28+
ocel.get_ob_types().collect::<Vec<_>>()
3429
);
30+
31+
let preview_n = 10;
32+
println!("First {} events:", preview_n);
33+
for ev in ocel.get_all_evs().take(preview_n) {
34+
let ev_type = ocel.get_ev_type_of(ev);
35+
let timestamp = ocel.get_ev_time(ev);
36+
println!(
37+
"Event {:?}: Type: {}, Timestamp: {}",
38+
ev, ev_type, timestamp
39+
);
40+
let attrs = ocel.get_ev_attrs(ev);
41+
for attr in attrs {
42+
let val = ocel.get_ev_attr_val(ev, attr);
43+
println!(" Attribute: {} = {:?}", attr, val);
44+
}
45+
}
3546
Ok(())
3647
}
Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
11
//! Object-centric Process Analysis
22
33
pub mod object_attribute_changes;
4+
pub mod oc_performance;
5+
pub mod oc_statistics;
Lines changed: 152 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,152 @@
1+
//! Object-centric performance analysis over [`SlimLinkedOCEL`]: per-event sojourn and
2+
//! synchronization times.
3+
4+
use macros_process_mining::register_binding;
5+
use rayon::prelude::*;
6+
7+
use crate::core::event_data::object_centric::linked_ocel::{
8+
slim_linked_ocel::{EventIndex, ObjectIndex},
9+
LinkedOCELAccess, SlimLinkedOCEL,
10+
};
11+
12+
/// Each object's reverse-E2O events in `(time, id)` order.
13+
///
14+
/// Events are sorted per call, so this makes no assumption about global event ordering
15+
fn sorted_events_per_object(ocel: &SlimLinkedOCEL) -> Vec<Vec<EventIndex>> {
16+
(0..ocel.get_num_obs() as u32)
17+
.into_par_iter()
18+
.map(|i| {
19+
let mut evs: Vec<EventIndex> =
20+
ObjectIndex::from(i).get_e2o_rev(ocel).copied().collect();
21+
evs.sort_by(|a, b| {
22+
a.get_time(ocel)
23+
.cmp(b.get_time(ocel))
24+
.then_with(|| ocel.get_ev_id(a).cmp(ocel.get_ev_id(b)))
25+
});
26+
evs
27+
})
28+
.collect()
29+
}
30+
31+
/// The `(time, id)`-immediate predecessor of `e` on object `o`, using the
32+
/// per-object sorted lists from [`sorted_events_per_object`].
33+
/// `None` if `e` is the first event on `o`.
34+
#[inline]
35+
fn df_predecessor(
36+
sorted: &[Vec<EventIndex>],
37+
ocel: &SlimLinkedOCEL,
38+
e: EventIndex,
39+
o: ObjectIndex,
40+
) -> Option<EventIndex> {
41+
let evs = &sorted[o.into_inner() as usize];
42+
let key = (e.get_time(ocel), ocel.get_ev_id(&e));
43+
let pos = evs
44+
.binary_search_by(|x| (x.get_time(ocel), ocel.get_ev_id(x)).cmp(&key))
45+
.ok()?;
46+
pos.checked_sub(1).map(|p| evs[p])
47+
}
48+
49+
/// Per-event synchronization time and the delaying object.
50+
///
51+
/// For each event with at least one directly-follows predecessor, the synchronization time is
52+
/// `max_predecessor_time - min_predecessor_time` in integer microseconds (the span between its
53+
/// earliest and latest directly-preceding event). The delaying object is the object linking the
54+
/// latest predecessor (ties broken by ascending object id).
55+
/// Returns one row `(event_id, sync_us, delaying_object_id)` per qualifying event.
56+
///
57+
/// `top_k`: if `Some(k)`, return only the `k` rows with the largest `sync_us`, ties broken by
58+
/// ascending event id, sorted descending. `None` returns every qualifying event.
59+
#[register_binding]
60+
pub fn locel_oc_perf_sync_per_event(
61+
ocel: &SlimLinkedOCEL,
62+
#[bind(default)] top_k: Option<usize>,
63+
) -> Vec<(String, i64, String)> {
64+
let sorted = sorted_events_per_object(ocel);
65+
let mut rows: Vec<(EventIndex, i64, ObjectIndex)> = (0..ocel.get_num_evs() as u32)
66+
.into_par_iter()
67+
.filter_map(|i| {
68+
let e = EventIndex::from(i);
69+
let mut min_us = i64::MAX;
70+
// (latest predecessor time, its object) = the delaying edge.
71+
let mut delaying: Option<(i64, ObjectIndex)> = None;
72+
for &o in e.get_e2o(ocel) {
73+
if let Some(p) = df_predecessor(&sorted, ocel, e, o) {
74+
let t = p.get_time(ocel).timestamp_micros();
75+
min_us = min_us.min(t);
76+
let keep = match delaying {
77+
Some((bt, bo)) => {
78+
bt > t || (bt == t && ocel.get_ob_id(&bo) <= ocel.get_ob_id(&o))
79+
}
80+
None => false,
81+
};
82+
if !keep {
83+
delaying = Some((t, o));
84+
}
85+
}
86+
}
87+
delaying.map(|(max_us, o)| (e, max_us - min_us, o))
88+
})
89+
.collect();
90+
if let Some(k) = top_k {
91+
let cmp = |a: &(EventIndex, i64, ObjectIndex), b: &(EventIndex, i64, ObjectIndex)| {
92+
b.1.cmp(&a.1)
93+
.then_with(|| ocel.get_ev_id(&a.0).cmp(ocel.get_ev_id(&b.0)))
94+
};
95+
if k < rows.len() {
96+
rows.select_nth_unstable_by(k, cmp);
97+
rows.truncate(k);
98+
}
99+
rows.sort_unstable_by(cmp);
100+
}
101+
rows.into_iter()
102+
.map(|(e, max_minus_min, o)| {
103+
(
104+
ocel.get_ev_id(&e).to_string(),
105+
max_minus_min,
106+
ocel.get_ob_id(&o).to_string(),
107+
)
108+
})
109+
.collect()
110+
}
111+
112+
/// Per-event sojourn time.
113+
///
114+
/// For each event with at least one directly-follows predecessor, the sojourn time is
115+
/// `event_time - latest_predecessor_time` in integer microseconds. Returns one row
116+
/// `(event_id, sojourn_us)` per qualifying event.
117+
///
118+
/// `top_k`: if `Some(k)`, return only the `k` rows with the largest `sojourn_us`, ties broken by
119+
/// ascending event id, sorted descending. `None` returns every qualifying event.
120+
#[register_binding]
121+
pub fn locel_oc_perf_sojourn_per_event(
122+
ocel: &SlimLinkedOCEL,
123+
#[bind(default)] top_k: Option<usize>,
124+
) -> Vec<(String, i64)> {
125+
let sorted = sorted_events_per_object(ocel);
126+
let mut rows: Vec<(EventIndex, i64)> = (0..ocel.get_num_evs() as u32)
127+
.into_par_iter()
128+
.filter_map(|i| {
129+
let e = EventIndex::from(i);
130+
let latest = e
131+
.get_e2o(ocel)
132+
.filter_map(|&o| df_predecessor(&sorted, ocel, e, o))
133+
.map(|p| p.get_time(ocel).timestamp_micros())
134+
.max()?;
135+
Some((e, e.get_time(ocel).timestamp_micros() - latest))
136+
})
137+
.collect();
138+
if let Some(k) = top_k {
139+
let cmp = |a: &(EventIndex, i64), b: &(EventIndex, i64)| {
140+
b.1.cmp(&a.1)
141+
.then_with(|| ocel.get_ev_id(&a.0).cmp(ocel.get_ev_id(&b.0)))
142+
};
143+
if k < rows.len() {
144+
rows.select_nth_unstable_by(k, cmp);
145+
rows.truncate(k);
146+
}
147+
rows.sort_unstable_by(cmp);
148+
}
149+
rows.into_iter()
150+
.map(|(e, sojourn_us)| (ocel.get_ev_id(&e).to_string(), sojourn_us))
151+
.collect()
152+
}
Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,79 @@
1+
//! Descriptive statistics over object-centric event data (E2O type counts, conversion rates).
2+
3+
use std::collections::HashMap;
4+
5+
use macros_process_mining::register_binding;
6+
use rayon::prelude::*;
7+
8+
use crate::core::event_data::object_centric::linked_ocel::{
9+
slim_linked_ocel::{EventIndex, ObjectIndex},
10+
LinkedOCELAccess, SlimLinkedOCEL,
11+
};
12+
13+
/// Count E2O relationships per `(event_type, object_type)` pair.
14+
///
15+
/// Each entry is `(event_type, object_type, count)`. A single `(event, object)` pair connected
16+
/// by multiple qualifiers contributes once per qualifier. Pairs with zero relations are omitted.
17+
/// Result row order is unspecified.
18+
#[register_binding]
19+
pub fn locel_event_object_type_counts(ocel: &SlimLinkedOCEL) -> Vec<(String, String, i64)> {
20+
let num_events = ocel.get_num_evs() as u32;
21+
let counts: HashMap<(usize, usize), i64> = (0..num_events)
22+
.into_par_iter()
23+
.fold(HashMap::new, |mut acc, i| {
24+
let ev = EventIndex::from(i).get_ev(ocel);
25+
for (_q, ob) in &ev.relationships {
26+
let ot = ob.get_ob(ocel).object_type;
27+
*acc.entry((ev.event_type, ot)).or_insert(0) += 1;
28+
}
29+
acc
30+
})
31+
.reduce(HashMap::new, merge_sum_maps);
32+
let ev_types: Vec<&str> = <SlimLinkedOCEL as LinkedOCELAccess>::get_ev_types(ocel).collect();
33+
let ob_types: Vec<&str> = <SlimLinkedOCEL as LinkedOCELAccess>::get_ob_types(ocel).collect();
34+
counts
35+
.into_iter()
36+
.map(|((e, o), c)| (ev_types[e].to_string(), ob_types[o].to_string(), c))
37+
.collect()
38+
}
39+
40+
/// Conversion rate from `source_type` to `target_type` via O2O, restricted to targets touched by `activity`.
41+
///
42+
/// Returns the fraction of `source_type` objects that have at least one outgoing O2O edge to a
43+
/// `target_type` object related (via E2O) to some event of the given event type. Returns `0.0`
44+
/// if no `source_type` objects exist.
45+
#[register_binding]
46+
pub fn locel_conversion_rate(
47+
ocel: &SlimLinkedOCEL,
48+
activity: String,
49+
source_type: String,
50+
target_type: String,
51+
) -> f64 {
52+
let sources: Vec<ObjectIndex> = ocel.get_obs_of_type(&source_type).copied().collect();
53+
let total = sources.len();
54+
if total == 0 {
55+
return 0.0;
56+
}
57+
let reached = sources
58+
.par_iter()
59+
.filter(|&&s| {
60+
s.get_o2o(ocel).any(|&t| {
61+
t.get_ob_type(ocel) == &target_type
62+
&& t.get_e2o_rev(ocel)
63+
.any(|&e| e.get_ev_type(ocel) == &activity)
64+
})
65+
})
66+
.count();
67+
reached as f64 / total as f64
68+
}
69+
70+
/// Merge `b` into `a` by summing `i64` counts for matching keys. Used as the rayon reduce step.
71+
fn merge_sum_maps<K: std::hash::Hash + Eq>(
72+
mut a: HashMap<K, i64>,
73+
b: HashMap<K, i64>,
74+
) -> HashMap<K, i64> {
75+
for (k, v) in b {
76+
*a.entry(k).or_insert(0) += v;
77+
}
78+
a
79+
}

0 commit comments

Comments
 (0)