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

Method sink

src/sinks/csv_file.rs:25–85  ·  view source on GitHub ↗
(
        &self,
        stream: &Stream<S, ResultDiff<u64>>,
        pact: P,
    )

Source from the content-addressed store, hash-verified

23
24impl Sinkable<u64> for CsvFile {
25 fn sink<S, P>(
26 &self,
27 stream: &Stream<S, ResultDiff<u64>>,
28 pact: P,
29 ) -> Result<Stream<S, ResultDiff<u64>>, Error>
30 where
31 S: Scope<Timestamp = u64>,
32 P: ParallelizationContract<S::Timestamp, ResultDiff<u64>>,
33 {
34 let writer_result = csv::WriterBuilder::new()
35 .has_headers(self.has_headers)
36 .delimiter(self.delimiter)
37 .from_path(&self.path);
38
39 match writer_result {
40 Err(error) => Err(Error::fault(format!("Failed to create writer: {}", error))),
41 Ok(mut writer) => {
42 let mut recvd = Vec::new();
43 let mut vector = Vec::new();
44
45 let name = format!("CsvFile({})", &self.path);
46
47 let mut builder = OperatorBuilder::new(name, stream.scope());
48 let mut input = builder.new_input(stream, pact);
49 let (_output, sunk) = builder.new_output();
50
51 builder.build(|_capabilities| {
52 move |frontiers| {
53 let mut input_handle =
54 FrontieredInputHandle::new(&mut input, &frontiers[0]);
55
56 input_handle.for_each(|_cap, data| {
57 data.swap(&mut vector);
58 // @TODO what to do with diff here?
59 for (tuple, time, _diff) in vector.drain(..) {
60 recvd.push((time, tuple));
61 }
62 });
63
64 recvd.sort_by(|x, y| x.0.cmp(&y.0));
65
66 // determine how many (which) elements to read from `recvd`.
67 let count = recvd
68 .iter()
69 .filter(|&(ref time, _)| !input_handle.frontier().less_equal(time))
70 .count();
71
72 for (_, tuple) in recvd.drain(..count) {
73 writer.serialize(tuple).expect("failed to write record");
74 }
75
76 if input_handle.frontier.is_empty() {
77 println!("Inputs to csv sink have ceased.");
78 }
79 }
80 });
81
82 Ok(sunk)

Callers

nothing calls this directly

Calls 3

cmpMethod · 0.80
countMethod · 0.80
less_equalMethod · 0.45

Tested by

no test coverage detected