(
&self,
stream: &Stream<S, ResultDiff<u64>>,
pact: P,
)
| 23 | |
| 24 | impl 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) |
nothing calls this directly
no test coverage detected