()
| 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 | // Kafka Config |
| 14 | let brokers = "localhost:9092"; |
| 15 | let topic_name = "topic1"; |
| 16 | let message_timeout_ms = "5000"; |
| 17 | |
| 18 | |
| 19 | let kafka_proc_factory = || KafkaProcessor::new(brokers, topic_name, message_timeout_ms); |
| 20 | let kafka_proc_concurrency = 1; |
| 21 | let kafka_proc_router = RouterType::RoundRobin; |
| 22 | let kafka_proc_buffer_size = 100; |
| 23 | |
| 24 | |
| 25 | let result_handler_proc_factory = || DeliveryHandler; |
| 26 | let result_handler_proc_concurrency = 1; |
| 27 | let result_handler_proc_buffer_size = 10; |
| 28 | |
| 29 | |
| 30 | |
| 31 | // 1. create X processor instances by 'proc_concurrency' |
| 32 | // |
| 33 | // 2. create X producer instances by 'producer_concurrency' |
| 34 | // |
| 35 | // 3. create topology and syncing |
| 36 | // |
| 37 | |
| 38 | |
| 39 | // / \ |
| 40 | // producer-1 / \ |
| 41 | // producer-2 ----- > Kafka_processor -> processor |
| 42 | // producer-3 \ / |
| 43 | // \ / |
| 44 | // |
| 45 | |
| 46 | let safe_shutdown = |
| 47 | run_topology_2( |
| 48 | producer_factory, |
| 49 | producer_concurrency, |
| 50 | producer_router, |
| 51 | producer_buffer_pool, |
| 52 | |
| 53 | kafka_proc_factory, |
| 54 | kafka_proc_concurrency, |
| 55 | kafka_proc_router, |
| 56 | kafka_proc_buffer_size, |
| 57 | |
| 58 | result_handler_proc_factory, |
| 59 | result_handler_proc_concurrency, |
| 60 | result_handler_proc_buffer_size |
| 61 | ); |
| 62 |
nothing calls this directly
no test coverage detected