()
| 1 | |
| 2 | #[tokio::main] |
| 3 | async fn main() { |
| 4 | |
| 5 | let producer_factory = || Prod; |
| 6 | let producer_concurrency = 3; |
| 7 | let producer_router = RouterType::RoundRobin; |
| 8 | let producer_buffer_pool = 100; |
| 9 | |
| 10 | |
| 11 | // Kafka Config |
| 12 | let brokers = "localhost:9092"; |
| 13 | let topic_name = "topic1"; |
| 14 | let message_timeout_ms = "5000"; |
| 15 | |
| 16 | let proc_factory = |
| 17 | || KafkaProcessor::new(brokers, |
| 18 | topic_name, |
| 19 | message_timeout_ms); |
| 20 | |
| 21 | let proc_concurrency = 1; |
| 22 | let proc_buffer_size = 100; |
| 23 | |
| 24 | |
| 25 | |
| 26 | // 1. create X processor instances by 'proc_concurrency' |
| 27 | // |
| 28 | // 2. create X producer instances by 'producer_concurrency' |
| 29 | // |
| 30 | // 3. create topology and syncing |
| 31 | // |
| 32 | |
| 33 | |
| 34 | // / \ |
| 35 | // producer-1 / \ |
| 36 | // producer-2 ----- > Kafka_processor |
| 37 | // producer-3 \ / |
| 38 | // \ / |
| 39 | // |
| 40 | |
| 41 | let safe_shutdown = |
| 42 | run_topology_1( |
| 43 | producer_factory, |
| 44 | producer_concurrency, |
| 45 | producer_router, |
| 46 | producer_buffer_pool, |
| 47 | |
| 48 | proc_factory, |
| 49 | proc_concurrency, |
| 50 | proc_buffer_size, |
| 51 | ); |
| 52 | |
| 53 | |
| 54 | // Safe Shutdown from (Producer) to (Layer_X_Processor) |
| 55 | safe_shutdown.send(()); |
| 56 | |
| 57 | } |
| 58 | |
| 59 | |
| 60 |
nothing calls this directly
no test coverage detected