use anyhow::Result;
use async_trait::async_trait;
use flume::Sender;
use std::time::Duration;
use tracing::{Level, info};
use work_dispatcher::{Processor, Producer, WorkDispatcher};
#[derive(Debug)]
struct CsvRow {
id: i32,
data: String,
}
#[derive(Clone)]
struct DbConnection(String);
struct CsvProducer {
count: usize,
}
#[async_trait]
impl Producer for CsvProducer {
type Item = CsvRow;
async fn run(self, sender: Sender<Self::Item>) {
info!("[Producer] Starting to generate {} records.", self.count);
for i in 0..self.count {
let row = CsvRow {
id: i as i32,
data: format!("Data for row {}", i),
};
sender.send_async(row).await.unwrap();
}
}
}
#[derive(Clone)]
struct DbProcessor;
#[async_trait]
impl Processor for DbProcessor {
type Item = CsvRow;
type Context = DbConnection;
async fn process(&self, item: Self::Item, context: &Self::Context) {
info!(
"[Processor] Writing item #{} data:{} to DB '{}'",
item.id, item.data, context.0
);
tokio::time::sleep(Duration::from_secs(1)).await;
}
}
#[tokio::main]
async fn main() -> Result<()> {
tracing_subscriber::fmt()
.with_max_level(Level::DEBUG)
.init();
let producer = CsvProducer { count: 100 };
let processor = DbProcessor;
let context = DbConnection("mongodb://localhost:27017".to_string());
WorkDispatcher::new(producer, processor, context)
.workers(8)
.buffer(500)
.run()
.await?;
Ok(())
}