1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
use serde::Serialize;
use std::thread::{spawn, JoinHandle};
use std::sync::mpsc::{channel, Sender};
use bson;
use mongodb::ThreadedClient;
use mongodb::db::ThreadedDatabase;
use super::Sink;
#[derive(new, Debug, Clone)]
pub struct MongoSink {
host: String,
port: u16,
db: String,
collection: String,
}
impl MongoSink {
pub fn local(db: &str, collection: &str) -> Self {
Self::new(
"localhost".to_string(),
27017,
db.to_string(),
collection.to_string(),
)
}
}
fn into_bson_document<T: Serialize>(val: T) -> bson::Document {
use bson::Bson::*;
match bson::to_bson(&val).unwrap() {
Document(d) => d,
_ => panic!("Input data must be converted into BSON::Document"),
}
}
impl<Doc: 'static + Send + Serialize> Sink<Doc> for MongoSink {
fn run(self) -> (Sender<Doc>, JoinHandle<()>) {
let (s, r) = channel::<Doc>();
let th = spawn(move || {
let cli = ::mongodb::Client::connect(&self.host, self.port)
.expect("Unable to connect to MongoDB");
let coll = cli.db(&self.db).collection(&self.collection);
loop {
match r.recv() {
Ok(doc) => {
coll.insert_one(into_bson_document(doc), None)
.expect("Failed to insert document");
}
Err(_) => break,
}
}
});
(s, th)
}
}