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