()
| 4 | |
| 5 | #[tokio::main] |
| 6 | async fn main() { |
| 7 | |
| 8 | let producer_factory = || Prod; |
| 9 | let producer_concurrency = 3; |
| 10 | let producer_router = RouterType::RoundRobin; |
| 11 | let producer_buffer_pool = 100; |
| 12 | |
| 13 | |
| 14 | // Pulsar Config |
| 15 | let addr = "pulsar://127.0.0.1:6650"; |
| 16 | let topic = "non-persistent://public/default/topic1"; |
| 17 | let pulsar_instance_name = "pulsar_processor"; |
| 18 | |
| 19 | let pulsar: Pulsar<_> = Pulsar::builder(addr, TokioExecutor).build().await.unwrap(); |
| 20 | let opts = producer::ProducerOptions { |
| 21 | schema: Some(proto::Schema { |
| 22 | r#type: proto::schema::Type::String as i32, |
| 23 | ..Default::default() |
| 24 | }), |
| 25 | ..Default::default() |
| 26 | }; |
| 27 | |
| 28 | |
| 29 | let pulsar_proc_factory = || || PulsarProcessor::new(pulsar, opts, topic, pulsar_instance_name).unwrap(); |
| 30 | let pulsar_proc_concurrency = 1; |
| 31 | let pulsar_proc_router = RouterType::RoundRobin; |
| 32 | let pulsar_proc_buffer_size = 100; |
| 33 | |
| 34 | |
| 35 | let result_handler_proc_factory = || DeliveryHandler; |
| 36 | let result_handler_proc_concurrency = 1; |
| 37 | let result_handler_proc_buffer_size = 10; |
| 38 | |
| 39 | |
| 40 | |
| 41 | // 1. create X processor instances by 'proc_concurrency' |
| 42 | // |
| 43 | // 2. create X producer instances by 'producer_concurrency' |
| 44 | // |
| 45 | // 3. create topology and syncing |
| 46 | // |
| 47 | |
| 48 | |
| 49 | // / \ |
| 50 | // producer-1 / \ |
| 51 | // producer-2 ----- > Pulsar_processor -> processor |
| 52 | // producer-3 \ / |
| 53 | // \ / |
| 54 | // |
| 55 | |
| 56 | let safe_shutdown = |
| 57 | run_topology_2( |
| 58 | producer_factory, |
| 59 | producer_concurrency, |
| 60 | producer_router, |
| 61 | producer_buffer_pool, |
| 62 | |
| 63 | pulsar_proc_factory, |
nothing calls this directly
no test coverage detected