From c8dab5d7ade4e958ab1c0837658055d182fb8be8 Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Mon, 26 Jan 2026 04:56:21 +0300 Subject: [PATCH] feat: run multiple tap consumers in parallel --- src/main.rs | 22 +++++++++++++++------- 1 file changed, 15 insertions(+), 7 deletions(-) diff --git a/src/main.rs b/src/main.rs index 5d5df06..8421023 100644 --- a/src/main.rs +++ b/src/main.rs @@ -34,13 +34,21 @@ async fn main() -> anyhow::Result<()> { let ops_count = Arc::new(AtomicU64::new(0)); - // start tap consumer - let db_clone = state.db.clone(); - let counts_clone = state.counts.clone(); - let ops_count_clone = ops_count.clone(); - tokio::spawn(async move { - run_tap_consumer(db_clone, counts_clone, ops_count_clone).await; - }); + // start tap consumers + let num_consumers = std::env::var("TAP_CONCURRENCY") + .ok() + .and_then(|s| s.parse().ok()) + .unwrap_or(100); + + for i in 0..num_consumers { + let db_clone = state.db.clone(); + let counts_clone = state.counts.clone(); + let ops_count_clone = ops_count.clone(); + tokio::spawn(async move { + info!("starting consumer #{}", i); + run_tap_consumer(db_clone, counts_clone, ops_count_clone).await; + }); + } // start stats reporter let ops_count_stats = ops_count.clone(); -- 2.51.2