MCPcopy Create free account
hub / github.com/comnik/declarative-dataflow / source

Method source

src/sources/differential_logging.rs:28–111  ·  view source on GitHub ↗
(
        &self,
        scope: &mut S,
        context: SourcingContext<S::Timestamp>,
    )

Source from the content-addressed store, hash-verified

26
27impl<S: Scope<Timestamp = Duration>> Sourceable<S> for DifferentialLogging {
28 fn source(
29 &self,
30 scope: &mut S,
31 context: SourcingContext<S::Timestamp>,
32 ) -> Vec<(
33 Aid,
34 AttributeConfig,
35 Stream<S, ((Value, Value), Duration, isize)>,
36 )> {
37 let input = Some(context.differential_events).replay_into(scope);
38
39 let mut demux =
40 OperatorBuilder::new("Differential Logging Demux".to_string(), scope.clone());
41
42 let mut input = demux.new_input(&input, Pipeline);
43
44 let mut wrappers = HashMap::with_capacity(self.attributes.len());
45 let mut streams = HashMap::with_capacity(self.attributes.len());
46
47 for aid in self.attributes.iter() {
48 let (wrapper, stream) = demux.new_output();
49 wrappers.insert(aid.to_string(), wrapper);
50 streams.insert(aid.to_string(), stream);
51 }
52
53 let mut demux_buffer = Vec::new();
54 let num_interests = self.attributes.len();
55
56 demux.build(move |_capability| {
57 move |_frontiers| {
58 let mut handles = HashMap::with_capacity(num_interests);
59 for (aid, wrapper) in wrappers.iter_mut() {
60 handles.insert(aid.to_string(), wrapper.activate());
61 }
62
63 input.for_each(|time, data: RefOrMut<Vec<_>>| {
64 data.swap(&mut demux_buffer);
65
66 let mut sessions = HashMap::with_capacity(num_interests);
67 for (aid, handle) in handles.iter_mut() {
68 sessions.insert(aid.to_string(), handle.session(&time));
69 }
70
71 for (time, _worker, datum) in demux_buffer.drain(..) {
72 match datum {
73 DifferentialEvent::Batch(x) => {
74 let operator = Eid((x.operator as u64).into());
75 let length = Number(x.length as i64);
76
77 sessions
78 .get_mut("differential.event/size")
79 .map(|s| s.give(((operator, length), time, 1)));
80 }
81 DifferentialEvent::Merge(x) => {
82 trace!("[DIFFERENTIAL] {:?}", x);
83
84 if let Some(complete_size) = x.complete {
85 let operator = Eid((x.operator as u64).into());

Callers

nothing calls this directly

Calls 1

intoMethod · 0.80

Tested by

no test coverage detected