From 397a7df95a4e80457fbd7821499a3bf0a85c7cf3 Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Fri, 28 Aug 2026 11:57:23 -0700 Subject: [PATCH 1/7] ref(pull): cleanup wiring in pull_consumer.rs to use factories where possible. --- rust_snuba/src/pull_consumer.rs | 549 ++++++++++++++++++-------------- 1 file changed, 303 insertions(+), 246 deletions(-) diff --git a/rust_snuba/src/pull_consumer.rs b/rust_snuba/src/pull_consumer.rs index 1ec5121453..79072ef342 100644 --- a/rust_snuba/src/pull_consumer.rs +++ b/rust_snuba/src/pull_consumer.rs @@ -10,11 +10,13 @@ use sentry_arroyo::processing::stream::{ }; use sentry_arroyo::types::{Topic, TopicOrPartition}; -use crate::config::{self, ProcessorConfig}; +use crate::config::{self, BatchSizeCalculation, ProcessorConfig, TopicConfig}; use crate::logging::{setup_logging, setup_sentry}; use crate::metrics::statsd::create_dogstatsd_backend; use crate::processors::{get_cogs_label, get_processing_function, ProcessingFunctionType}; +use crate::pull::batch::batch_metadata::BatchMetadata; use crate::pull::batch::buffer::PipelineBatchBuffer; +use crate::pull::batch::pipeline_batch::PipelineBatch; use crate::pull::pipelines::eap::EapPipeline; use crate::pull::pipelines::fire_and_forget::FireAndForgetPipeline; use crate::pull::producers::DryRunProducer; @@ -39,6 +41,266 @@ const FIRE_AND_FORGET_PROCESSORS: &[&str] = &[ /// Allowed processors for the EAP pipeline. const EAP_PROCESSORS: &[&str] = &["EAPItemsProcessor"]; +// ── Atomic factories ─────────────────────────────────────── + +fn resolve_processor( + processor_name: &str, +) -> Result { + match get_processing_function(processor_name) { + Some(ProcessingFunctionType::ProcessingFunction(f)) => Ok(f), + Some(ProcessingFunctionType::ProcessingFunctionWithReplacements(_)) => Err(format!( + "{processor_name} is a replacement processor — not supported" + )), + None => Err(format!("Unknown processor: {processor_name}")), + } +} + +fn make_kafka_source( + consumer_config: &config::ConsumerConfig, + consumer_group: &str, + auto_offset_reset: &str, + no_strict_offset_reset: bool, + max_poll_interval_ms: usize, +) -> KafkaSource { + let kafka_config = KafkaConfig::new_consumer_config( + vec![], + consumer_group.to_owned(), + auto_offset_reset + .parse() + .expect("Invalid auto_offset_reset"), + !no_strict_offset_reset, + max_poll_interval_ms, + Some(consumer_config.raw_topic.broker_config.clone()), + ); + let topic = Topic::new(&consumer_config.raw_topic.physical_topic_name); + KafkaSource::new(kafka_config, &[topic]) +} + +fn make_processor_config( + storage: &config::StorageConfig, + env_config: &config::EnvConfig, +) -> ProcessorConfig { + ProcessorConfig { + env_config: env_config.clone(), + storage_name: storage.name.clone(), + eap_items_emit_received_at: crate::processors::eap_items::emit_received_at(), + } +} + +fn make_kafka_producer(topic_config: &TopicConfig) -> KafkaProducer { + KafkaProducer::new(KafkaConfig::new_producer_config( + vec![], + Some(topic_config.broker_config.clone()), + )) +} + +fn make_writer( + storage: &config::StorageConfig, + format: InsertFormat, + columns: Option<&'static [&'static str]>, + dry_run_latency: Option, +) -> ClickHouseWriterStage { + if let Some(latency) = dry_run_latency { + ClickHouseWriterStage::new(DryRunWriter::new(latency)) + } else { + ClickHouseWriterStage::new(ClickhouseClient::new( + &storage.clickhouse_cluster, + &storage.clickhouse_table_name, + storage.name.clone(), + format, + columns, + )) + } +} + +fn make_dlq(topic_config: Option<&TopicConfig>, dry_run: bool) -> DlqHandler { + match topic_config { + Some(tc) if !dry_run => DlqHandler::new( + make_kafka_producer(tc), + TopicOrPartition::Topic(Topic::new(&tc.physical_topic_name)), + ), + _ => DlqHandler::new( + DryRunProducer, + TopicOrPartition::Topic(Topic::new("dry-run-dlq")), + ), + } +} + +fn make_commit_log( + topic_config: Option<&TopicConfig>, + source_topic: &str, + consumer_group: &str, + dry_run: bool, +) -> CommitLogStage { + match topic_config { + Some(tc) if !dry_run => CommitLogStage::new( + make_kafka_producer(tc), + Topic::new(&tc.physical_topic_name), + Topic::new(source_topic), + consumer_group.to_string(), + ), + _ => CommitLogStage::new( + DryRunProducer, + Topic::new("dry-run-commit-log"), + Topic::new(source_topic), + consumer_group.to_string(), + ), + } +} + +fn make_cogs( + topic_config: &TopicConfig, + resource_id: String, + dry_run: bool, + record_cogs: bool, +) -> CogsStage { + if !dry_run && record_cogs { + CogsStage::new( + make_kafka_producer(topic_config), + Topic::new(&topic_config.physical_topic_name), + resource_id, + ) + } else { + CogsStage::new(DryRunProducer, Topic::new("dry-run-cogs"), resource_id) + } +} + +fn make_processor_stage( + processor: crate::processors::ProcessingFunction, + processor_config: &ProcessorConfig, +) -> ProcessorStage { + ProcessorStage::new(processor, processor_config.clone()) +} + +fn make_batch( + max_batch_size: u64, + calculation: BatchSizeCalculation, +) -> BatchStage { + let (max_rows, max_bytes) = match calculation { + BatchSizeCalculation::Rows => (max_batch_size, u64::MAX), + BatchSizeCalculation::Bytes => (u64::MAX, max_batch_size), + }; + BatchStage::new(PipelineBatchBuffer::new(), max_rows, max_bytes) +} + +/// Returns (cadence, idle_timeout) for apply_with_timer. +fn make_flush_timers( + consumer_config: &config::ConsumerConfig, +) -> (Option, Option) { + let cadence = Some(Duration::from_millis(consumer_config.max_batch_time_ms)); + let idle_timeout = None; + (cadence, idle_timeout) +} + +// ── Pipeline assembly ────────────────────────────────────── + +fn make_eap_pipeline( + consumer_config: &config::ConsumerConfig, + storage: &config::StorageConfig, + processor_config: &ProcessorConfig, + consumer_group: &str, + processing_concurrency: usize, + clickhouse_concurrency: usize, + dry_run: bool, + dry_run_latency: Option, +) -> EapPipeline { + let processor_name = &storage.message_processor.python_class_name; + let source_topic_name = &consumer_config.raw_topic.physical_topic_name; + + assert_eq!(processor_name, "EAPItemsProcessor"); + + let insert_columns = Some(crate::processors::eap_items::EAPItemRow::column_names( + processor_config.eap_items_emit_received_at, + )); + let resource_id = + get_cogs_label(processor_name).unwrap_or_else(|| format!("{}_processor", storage.name)); + let (cadence, idle_timeout) = make_flush_timers(consumer_config); + + EapPipeline::new( + make_processor_stage( + crate::processors::eap_items::process_message_row_binary, + processor_config, + ), + processing_concurrency, + make_dlq(consumer_config.dlq_topic.as_ref(), dry_run), + make_batch( + consumer_config.max_batch_size as u64, + consumer_config.max_batch_size_calculation, + ), + cadence, + idle_timeout, + make_writer( + storage, + InsertFormat::RowBinary, + insert_columns, + dry_run_latency, + ), + clickhouse_concurrency, + make_commit_log( + consumer_config.commit_log_topic.as_ref(), + source_topic_name, + consumer_group, + dry_run, + ), + make_cogs( + &consumer_config.accountant_topic, + resource_id, + dry_run, + consumer_config.env.record_cogs, + ), + ) +} + +fn make_faf_pipeline( + processor: crate::processors::ProcessingFunction, + consumer_config: &config::ConsumerConfig, + storage: &config::StorageConfig, + processor_config: &ProcessorConfig, + processing_concurrency: usize, + clickhouse_concurrency: usize, + dry_run_latency: Option, +) -> FireAndForgetPipeline { + let (cadence, idle_timeout) = make_flush_timers(consumer_config); + + FireAndForgetPipeline::new( + make_processor_stage(processor, processor_config), + processing_concurrency, + make_batch( + consumer_config.max_batch_size as u64, + consumer_config.max_batch_size_calculation, + ), + cadence, + idle_timeout, + make_writer(storage, InsertFormat::JsonEachRow, None, dry_run_latency), + clickhouse_concurrency, + ) +} + +// ── Rebalance loop ───────────────────────────────────────── + +async fn run_with_rebalance>( + source: &KafkaSource, + build_pipeline: impl Fn() -> P, +) -> usize { + loop { + let pipeline = build_pipeline(); + let mut tracker = OffsetTracker::new(Duration::from_secs(1), source.committer()); + match pipeline.stream(source.stream()).commit(&mut tracker).await { + Ok(PipelineExit::Rebalance) => { + tracing::info!("Rebalance detected, restarting pipeline"); + continue; + } + Ok(PipelineExit::Shutdown | PipelineExit::Complete) => return 0, + Err(e) => { + tracing::error!("Pipeline failed: {}", e); + return 1; + } + } + } +} + +// ── Entry point ──────────────────────────────────────────── + #[pyfunction] #[allow(clippy::too_many_arguments)] pub fn pull_consumer( @@ -115,15 +377,10 @@ fn pull_consumer_impl( } } - // Resolve processor - let processor = match get_processing_function(&processor_name) { - Some(ProcessingFunctionType::ProcessingFunction(f)) => f, - Some(ProcessingFunctionType::ProcessingFunctionWithReplacements(_)) => { - tracing::error!("{processor_name} is a replacement processor — not supported"); - return 1; - } - None => { - tracing::error!("Unknown processor: {processor_name}"); + let processor = match resolve_processor(&processor_name) { + Ok(f) => f, + Err(msg) => { + tracing::error!("{msg}"); return 1; } }; @@ -137,6 +394,11 @@ fn pull_consumer_impl( } let dry_run = dry_run_latency_ms > 0; + let dry_run_latency = if dry_run { + Some(Duration::from_millis(dry_run_latency_ms)) + } else { + None + }; tracing::info!( storage = storage.name, @@ -147,68 +409,47 @@ fn pull_consumer_impl( "Starting pull consumer", ); - // Kafka source - let kafka_config = KafkaConfig::new_consumer_config( - vec![], - consumer_group.to_owned(), - auto_offset_reset - .parse() - .expect("Invalid auto_offset_reset"), - !no_strict_offset_reset, - max_poll_interval_ms, - Some(consumer_config.raw_topic.broker_config.clone()), - ); - let topic = Topic::new(&consumer_config.raw_topic.physical_topic_name); - - let eap_items_emit_received_at = crate::processors::eap_items::emit_received_at(); - - let processor_config = ProcessorConfig { - env_config: env_config.clone(), - storage_name: storage.name.clone(), - eap_items_emit_received_at, - }; - - let max_batch_size = consumer_config.max_batch_size as u64; - let max_batch_time = Duration::from_millis(consumer_config.max_batch_time_ms); - - // Build tokio runtime and run let rt = tokio::runtime::Builder::new_multi_thread() .enable_all() .build() .expect("failed to build tokio runtime"); let exit_code = rt.block_on(async { - // KafkaSource must be created inside the tokio runtime - // (rdkafka's StreamConsumer requires a reactor) - let source = KafkaSource::new(kafka_config, &[topic]); + let source = make_kafka_source( + &consumer_config, + consumer_group, + auto_offset_reset, + no_strict_offset_reset, + max_poll_interval_ms, + ); + let processor_config = make_processor_config(&storage, &env_config); let result = if is_eap { - run_eap( - &source, - &processor_config, - &storage, - &consumer_config, - &env_config, - consumer_group, - processing_concurrency, - max_batch_size, - max_batch_time, - clickhouse_concurrency, - dry_run_latency_ms, - ) + run_with_rebalance(&source, || { + make_eap_pipeline( + &consumer_config, + &storage, + &processor_config, + consumer_group, + processing_concurrency, + clickhouse_concurrency, + dry_run, + dry_run_latency, + ) + }) .await } else { - run_fire_and_forget( - &source, - processor, - &processor_config, - &storage, - processing_concurrency, - max_batch_size, - max_batch_time, - clickhouse_concurrency, - dry_run_latency_ms, - ) + run_with_rebalance(&source, || { + make_faf_pipeline( + processor, + &consumer_config, + &storage, + &processor_config, + processing_concurrency, + clickhouse_concurrency, + dry_run_latency, + ) + }) .await }; @@ -218,187 +459,3 @@ fn pull_consumer_impl( exit_code } - -#[allow(clippy::too_many_arguments)] -async fn run_fire_and_forget( - source: &KafkaSource, - processor: crate::processors::ProcessingFunction, - processor_config: &ProcessorConfig, - storage: &config::StorageConfig, - processing_concurrency: usize, - max_batch_size: u64, - max_batch_time: Duration, - clickhouse_concurrency: usize, - dry_run_latency_ms: u64, -) -> usize { - let dry_run = dry_run_latency_ms > 0; - loop { - let writer = if dry_run { - ClickHouseWriterStage::new(DryRunWriter::new(Duration::from_millis(dry_run_latency_ms))) - } else { - ClickHouseWriterStage::new(ClickhouseClient::new( - &storage.clickhouse_cluster, - &storage.clickhouse_table_name, - storage.name.clone(), - InsertFormat::JsonEachRow, - None, - )) - }; - - let pipeline = FireAndForgetPipeline::new( - ProcessorStage::new(processor, processor_config.clone()), - processing_concurrency, - BatchStage::new(PipelineBatchBuffer::new(), max_batch_size, u64::MAX), - Some(max_batch_time), - None, - writer, - clickhouse_concurrency, - ); - - let mut tracker = OffsetTracker::new(Duration::from_secs(1), source.committer()); - match pipeline.stream(source.stream()).commit(&mut tracker).await { - Ok(PipelineExit::Rebalance) => { - tracing::info!("Rebalance detected, restarting pipeline"); - continue; - } - Ok(PipelineExit::Shutdown | PipelineExit::Complete) => return 0, - Err(e) => { - tracing::error!("Pipeline failed: {}", e); - return 1; - } - } - } -} - -#[allow(clippy::too_many_arguments)] -async fn run_eap( - source: &KafkaSource, - processor_config: &ProcessorConfig, - storage: &config::StorageConfig, - consumer_config: &config::ConsumerConfig, - env_config: &config::EnvConfig, - consumer_group: &str, - processing_concurrency: usize, - max_batch_size: u64, - max_batch_time: Duration, - clickhouse_concurrency: usize, - dry_run_latency_ms: u64, -) -> usize { - let dry_run = dry_run_latency_ms > 0; - let processor_name = &storage.message_processor.python_class_name; - let source_topic_name = &consumer_config.raw_topic.physical_topic_name; - - // EAP pipeline is EAPItemsProcessor-only; RowBinary matches factory_v2.rs. - assert_eq!(processor_name, "EAPItemsProcessor"); - let processor = crate::processors::eap_items::process_message_row_binary; - let insert_format = InsertFormat::RowBinary; - let insert_columns: Option<&'static [&'static str]> = - Some(crate::processors::eap_items::EAPItemRow::column_names( - processor_config.eap_items_emit_received_at, - )); - - loop { - let writer = if dry_run { - ClickHouseWriterStage::new(DryRunWriter::new(Duration::from_millis(dry_run_latency_ms))) - } else { - ClickHouseWriterStage::new(ClickhouseClient::new( - &storage.clickhouse_cluster, - &storage.clickhouse_table_name, - storage.name.clone(), - insert_format, - insert_columns, - )) - }; - - let dlq_handler = if dry_run { - DlqHandler::new( - DryRunProducer, - TopicOrPartition::Topic(Topic::new("dry-run-dlq")), - ) - } else if let Some(ref topic_config) = consumer_config.dlq_topic { - let producer = KafkaProducer::new(KafkaConfig::new_producer_config( - vec![], - Some(topic_config.broker_config.clone()), - )); - DlqHandler::new( - producer, - TopicOrPartition::Topic(Topic::new(&topic_config.physical_topic_name)), - ) - } else { - DlqHandler::new( - DryRunProducer, - TopicOrPartition::Topic(Topic::new("no-dlq-configured")), - ) - }; - - let commit_log = if dry_run { - CommitLogStage::new( - DryRunProducer, - Topic::new("dry-run-commit-log"), - Topic::new(source_topic_name), - consumer_group.to_string(), - ) - } else if let Some(ref topic_config) = consumer_config.commit_log_topic { - let producer = KafkaProducer::new(KafkaConfig::new_producer_config( - vec![], - Some(topic_config.broker_config.clone()), - )); - CommitLogStage::new( - producer, - Topic::new(&topic_config.physical_topic_name), - Topic::new(source_topic_name), - consumer_group.to_string(), - ) - } else { - CommitLogStage::new( - DryRunProducer, - Topic::new("no-commit-log-configured"), - Topic::new(source_topic_name), - consumer_group.to_string(), - ) - }; - - let resource_id = - get_cogs_label(processor_name).unwrap_or_else(|| format!("{}_processor", storage.name)); - - let cogs = if dry_run || !env_config.record_cogs { - CogsStage::new(DryRunProducer, Topic::new("dry-run-cogs"), resource_id) - } else { - let producer = KafkaProducer::new(KafkaConfig::new_producer_config( - vec![], - Some(consumer_config.accountant_topic.broker_config.clone()), - )); - CogsStage::new( - producer, - Topic::new(&consumer_config.accountant_topic.physical_topic_name), - resource_id, - ) - }; - - let pipeline = EapPipeline::new( - ProcessorStage::new(processor, processor_config.clone()), - processing_concurrency, - dlq_handler, - BatchStage::new(PipelineBatchBuffer::new(), max_batch_size, u64::MAX), - Some(max_batch_time), - None, - writer, - clickhouse_concurrency, - commit_log, - cogs, - ); - - let mut tracker = OffsetTracker::new(Duration::from_secs(1), source.committer()); - match pipeline.stream(source.stream()).commit(&mut tracker).await { - Ok(PipelineExit::Rebalance) => { - tracing::info!("Rebalance detected, restarting pipeline"); - continue; - } - Ok(PipelineExit::Shutdown | PipelineExit::Complete) => return 0, - Err(e) => { - tracing::error!("Pipeline failed: {}", e); - return 1; - } - } - } -} From 99a12915d79dd1bb60cdeface3ca7e9746770864 Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Fri, 28 Aug 2026 15:29:16 -0700 Subject: [PATCH 2/7] ref(pull): Factor out shared clients to AppState pattern. Wire up use-row-binary. --- rust_snuba/Cargo.lock | 1 - rust_snuba/Cargo.toml | 3 + .../pull/stages/clickhouse_writer_stage.rs | 9 +- rust_snuba/src/pull/stages/cogs_stage.rs | 4 +- .../src/pull/stages/commit_log_stage.rs | 4 +- rust_snuba/src/pull/test_fixtures/sources.rs | 3 +- .../src/pull/tests/eap_pipeline_tests.rs | 8 +- .../pull/tests/processor_pipeline_tests.rs | 2 +- .../src/pull/writer/clickhouse_writer.rs | 6 - rust_snuba/src/pull_consumer.rs | 271 ++++++++++++------ snuba/cli/rust_consumer.py | 1 + 11 files changed, 200 insertions(+), 112 deletions(-) diff --git a/rust_snuba/Cargo.lock b/rust_snuba/Cargo.lock index 8039e9b538..34386b24c4 100644 --- a/rust_snuba/Cargo.lock +++ b/rust_snuba/Cargo.lock @@ -3912,7 +3912,6 @@ dependencies = [ [[package]] name = "sentry_arroyo" version = "2.43.0" -source = "git+https://github.com/getsentry/arroyo?branch=rbroughan%2Fpull-based-arroyo-part-2#ad99a41a4a59f62a0b176871a46cab2e4402c0c6" dependencies = [ "async-stream", "chrono", diff --git a/rust_snuba/Cargo.toml b/rust_snuba/Cargo.toml index da627a2a9b..503ebae250 100644 --- a/rust_snuba/Cargo.toml +++ b/rust_snuba/Cargo.toml @@ -104,3 +104,6 @@ tikv-jemallocator = "0.5" [[bench]] name = "processors" harness = false + +[patch."https://github.com/getsentry/arroyo"] +sentry_arroyo = { path = "/Users/ryanbroughan/code/arroyo" } diff --git a/rust_snuba/src/pull/stages/clickhouse_writer_stage.rs b/rust_snuba/src/pull/stages/clickhouse_writer_stage.rs index 3730875392..8d5d5dac6e 100644 --- a/rust_snuba/src/pull/stages/clickhouse_writer_stage.rs +++ b/rust_snuba/src/pull/stages/clickhouse_writer_stage.rs @@ -1,3 +1,4 @@ +use std::sync::Arc; use std::time::Instant; use chrono::Utc; @@ -14,14 +15,12 @@ use crate::pull::writer::ClickHouseWriter; /// commit log offsets, COGS data, and write stats for downstream /// handlers. Row bytes are freed after the write. pub struct ClickHouseWriterStage { - writer: Box, + writer: Arc, } impl ClickHouseWriterStage { - pub fn new(writer: impl ClickHouseWriter + 'static) -> Self { - Self { - writer: Box::new(writer), - } + pub fn new(writer: Arc) -> Self { + Self { writer } } } diff --git a/rust_snuba/src/pull/stages/cogs_stage.rs b/rust_snuba/src/pull/stages/cogs_stage.rs index 945ab38787..cd24234f75 100644 --- a/rust_snuba/src/pull/stages/cogs_stage.rs +++ b/rust_snuba/src/pull/stages/cogs_stage.rs @@ -37,12 +37,12 @@ pub struct CogsStage { impl CogsStage { pub fn new( - producer: impl Producer + 'static, + producer: Arc>, destination: Topic, resource_id: String, ) -> Self { Self { - producer: Arc::new(producer), + producer, destination: TopicOrPartition::Topic(destination), resource_id, logged_warning: AtomicBool::new(false), diff --git a/rust_snuba/src/pull/stages/commit_log_stage.rs b/rust_snuba/src/pull/stages/commit_log_stage.rs index 415929f308..bb286d9953 100644 --- a/rust_snuba/src/pull/stages/commit_log_stage.rs +++ b/rust_snuba/src/pull/stages/commit_log_stage.rs @@ -31,13 +31,13 @@ pub struct CommitLogStage { impl CommitLogStage { pub fn new( - producer: impl Producer + 'static, + producer: Arc>, destination: Topic, source_topic: Topic, consumer_group: String, ) -> Self { Self { - producer: Arc::new(producer), + producer, destination: TopicOrPartition::Topic(destination), source_topic, consumer_group, diff --git a/rust_snuba/src/pull/test_fixtures/sources.rs b/rust_snuba/src/pull/test_fixtures/sources.rs index 3218d9544e..b65f02ac01 100644 --- a/rust_snuba/src/pull/test_fixtures/sources.rs +++ b/rust_snuba/src/pull/test_fixtures/sources.rs @@ -6,7 +6,6 @@ use sentry_arroyo::processing::stream::{ MessageMetadata, OffsetCommitter, PipelineEnvelope, PullSource, StageResult, }; use sentry_arroyo::types::{Partition, Topic}; -use std::sync::Arc; /// In-memory source for testing. Constructs envelopes from raw payloads /// with sequential offsets on partition 0. @@ -26,7 +25,7 @@ impl VecSource { offset: i as u64, timestamp: chrono::Utc::now(), }; - StageResult::Emit(PipelineEnvelope::new(kp.clone(), md, Arc::new(kp))) + StageResult::Emit(PipelineEnvelope::new(kp.clone(), md, kp)) }) .collect(); diff --git a/rust_snuba/src/pull/tests/eap_pipeline_tests.rs b/rust_snuba/src/pull/tests/eap_pipeline_tests.rs index 793516ed5f..59c0ddcc81 100644 --- a/rust_snuba/src/pull/tests/eap_pipeline_tests.rs +++ b/rust_snuba/src/pull/tests/eap_pipeline_tests.rs @@ -63,22 +63,22 @@ async fn test_eap_pipeline() { ProcessorStage::new(processor, ProcessorConfig::default()), 1, // processing_concurrency DlqHandler::new( - dlq_producer, + Arc::new(dlq_producer), TopicOrPartition::Topic(Topic::new("snuba-dead-letter-items")), ), BatchStage::new(PipelineBatchBuffer::new(), 2, u64::MAX), Some(Duration::from_secs(2)), None, - ClickHouseWriterStage::new(Arc::clone(&writer)), + ClickHouseWriterStage::new(writer.clone()), 2, CommitLogStage::new( - commit_log_producer, + Arc::new(commit_log_producer), Topic::new("snuba-items-commit-log"), Topic::new("snuba-items"), "test-group".to_string(), ), CogsStage::new( - cogs_producer, + Arc::new(cogs_producer), Topic::new("shared-resources-usage"), "eap_items_processor".to_string(), ), diff --git a/rust_snuba/src/pull/tests/processor_pipeline_tests.rs b/rust_snuba/src/pull/tests/processor_pipeline_tests.rs index 9d77490949..1a5e8058d4 100644 --- a/rust_snuba/src/pull/tests/processor_pipeline_tests.rs +++ b/rust_snuba/src/pull/tests/processor_pipeline_tests.rs @@ -72,7 +72,7 @@ async fn test_fire_and_forget_pipeline(#[case] processor_name: &str, #[case] top BatchStage::new(PipelineBatchBuffer::new(), 2, u64::MAX), Some(Duration::from_secs(2)), None, - ClickHouseWriterStage::new(Arc::clone(&writer)), + ClickHouseWriterStage::new(writer.clone()), 2, ); diff --git a/rust_snuba/src/pull/writer/clickhouse_writer.rs b/rust_snuba/src/pull/writer/clickhouse_writer.rs index e83e5502e5..583fb2e4e9 100644 --- a/rust_snuba/src/pull/writer/clickhouse_writer.rs +++ b/rust_snuba/src/pull/writer/clickhouse_writer.rs @@ -13,9 +13,3 @@ use futures::future::BoxFuture; pub trait ClickHouseWriter: Send + Sync { fn write(&self, body: Vec) -> BoxFuture<'_, anyhow::Result<()>>; } - -impl ClickHouseWriter for std::sync::Arc { - fn write(&self, body: Vec) -> BoxFuture<'_, anyhow::Result<()>> { - (**self).write(body) - } -} diff --git a/rust_snuba/src/pull_consumer.rs b/rust_snuba/src/pull_consumer.rs index 79072ef342..dfa6254244 100644 --- a/rust_snuba/src/pull_consumer.rs +++ b/rust_snuba/src/pull_consumer.rs @@ -1,8 +1,11 @@ +use std::sync::Arc; use std::time::Duration; use pyo3::prelude::*; use sentry_arroyo::backends::kafka::config::KafkaConfig; use sentry_arroyo::backends::kafka::producer::KafkaProducer; +use sentry_arroyo::backends::kafka::types::KafkaPayload; +use sentry_arroyo::backends::Producer; use sentry_arroyo::metrics; use sentry_arroyo::processing::stream::{ BatchStage, DlqHandler, KafkaSource, OffsetTracker, Pipeline, PipelineExit, PipelineExt, @@ -13,6 +16,7 @@ use sentry_arroyo::types::{Topic, TopicOrPartition}; use crate::config::{self, BatchSizeCalculation, ProcessorConfig, TopicConfig}; use crate::logging::{setup_logging, setup_sentry}; use crate::metrics::statsd::create_dogstatsd_backend; +use crate::processors::eap_items::EAPItemRow; use crate::processors::{get_cogs_label, get_processing_function, ProcessingFunctionType}; use crate::pull::batch::batch_metadata::BatchMetadata; use crate::pull::batch::buffer::PipelineBatchBuffer; @@ -24,7 +28,7 @@ use crate::pull::stages::clickhouse_writer_stage::ClickHouseWriterStage; use crate::pull::stages::cogs_stage::CogsStage; use crate::pull::stages::commit_log_stage::CommitLogStage; use crate::pull::stages::processor_stage::ProcessorStage; -use crate::pull::writer::DryRunWriter; +use crate::pull::writer::{ClickHouseWriter, DryRunWriter}; use crate::strategies::clickhouse::writer_v2::{ClickhouseClient, InsertFormat}; /// Allowed processors for the fire-and-forget pipeline. @@ -41,7 +45,108 @@ const FIRE_AND_FORGET_PROCESSORS: &[&str] = &[ /// Allowed processors for the EAP pipeline. const EAP_PROCESSORS: &[&str] = &["EAPItemsProcessor"]; -// ── Atomic factories ─────────────────────────────────────── +// ── Shared resources ────────────────────────────────────── + +struct SharedResources { + dlq_producer: Arc>, + commit_log_producer: Arc>, + cogs_producer: Arc>, + ch_writer: Arc, +} + +fn make_shared_resources( + consumer_config: &config::ConsumerConfig, + storage: &config::StorageConfig, + processor_config: &ProcessorConfig, + use_row_binary: bool, + dry_run: bool, + dry_run_latency: Option, +) -> SharedResources { + SharedResources { + dlq_producer: Arc::from(make_dlq_producer(consumer_config, dry_run)), + commit_log_producer: Arc::from(make_commit_log_producer(consumer_config, dry_run)), + cogs_producer: Arc::from(make_cogs_producer(consumer_config, dry_run)), + ch_writer: Arc::from(make_ch_writer( + storage, + processor_config, + use_row_binary, + dry_run_latency, + )), + } +} + +// ── Producer factories ──────────────────── + +fn make_kafka_producer(topic_config: &TopicConfig) -> KafkaProducer { + KafkaProducer::new(KafkaConfig::new_producer_config( + vec![], + Some(topic_config.broker_config.clone()), + )) +} + +fn make_dlq_producer( + consumer_config: &config::ConsumerConfig, + dry_run: bool, +) -> Box> { + match consumer_config.dlq_topic.as_ref() { + Some(tc) if !dry_run => Box::new(make_kafka_producer(tc)), + _ => Box::new(DryRunProducer), + } +} + +fn make_commit_log_producer( + consumer_config: &config::ConsumerConfig, + dry_run: bool, +) -> Box> { + match consumer_config.commit_log_topic.as_ref() { + Some(tc) if !dry_run => Box::new(make_kafka_producer(tc)), + _ => Box::new(DryRunProducer), + } +} + +fn make_cogs_producer( + consumer_config: &config::ConsumerConfig, + dry_run: bool, +) -> Box> { + if !dry_run && consumer_config.env.record_cogs { + Box::new(make_kafka_producer(&consumer_config.accountant_topic)) + } else { + Box::new(DryRunProducer) + } +} + +fn make_ch_writer( + storage: &config::StorageConfig, + processor_config: &ProcessorConfig, + use_row_binary: bool, + dry_run_latency: Option, +) -> Box { + if let Some(latency) = dry_run_latency { + return Box::new(DryRunWriter::new(latency)); + } + + if use_row_binary { + Box::new(ClickhouseClient::new( + &storage.clickhouse_cluster, + &storage.clickhouse_table_name, + storage.name.clone(), + InsertFormat::RowBinary, + Some(EAPItemRow::column_names( + processor_config.eap_items_emit_received_at, + )), + )) + } else { + Box::new(ClickhouseClient::new( + &storage.clickhouse_cluster, + &storage.clickhouse_table_name, + storage.name.clone(), + InsertFormat::JsonEachRow, + None, + )) + } +} + +// ── Stage factories ─────────────────────────────────────── fn resolve_processor( processor_name: &str, @@ -87,82 +192,57 @@ fn make_processor_config( } } -fn make_kafka_producer(topic_config: &TopicConfig) -> KafkaProducer { - KafkaProducer::new(KafkaConfig::new_producer_config( - vec![], - Some(topic_config.broker_config.clone()), - )) -} - -fn make_writer( - storage: &config::StorageConfig, - format: InsertFormat, - columns: Option<&'static [&'static str]>, - dry_run_latency: Option, -) -> ClickHouseWriterStage { - if let Some(latency) = dry_run_latency { - ClickHouseWriterStage::new(DryRunWriter::new(latency)) - } else { - ClickHouseWriterStage::new(ClickhouseClient::new( - &storage.clickhouse_cluster, - &storage.clickhouse_table_name, - storage.name.clone(), - format, - columns, - )) - } -} - -fn make_dlq(topic_config: Option<&TopicConfig>, dry_run: bool) -> DlqHandler { - match topic_config { - Some(tc) if !dry_run => DlqHandler::new( - make_kafka_producer(tc), - TopicOrPartition::Topic(Topic::new(&tc.physical_topic_name)), - ), - _ => DlqHandler::new( - DryRunProducer, - TopicOrPartition::Topic(Topic::new("dry-run-dlq")), - ), - } +fn make_dlq_handler( + producer: &Arc>, + topic_config: Option<&TopicConfig>, + dry_run: bool, +) -> DlqHandler { + let topic_name = match topic_config { + Some(tc) if !dry_run => tc.physical_topic_name.as_str(), + _ => "dry-run-dlq", + }; + DlqHandler::new( + Arc::clone(producer), + TopicOrPartition::Topic(Topic::new(topic_name)), + ) } -fn make_commit_log( +fn make_commit_log_stage( + producer: &Arc>, topic_config: Option<&TopicConfig>, source_topic: &str, consumer_group: &str, dry_run: bool, ) -> CommitLogStage { - match topic_config { - Some(tc) if !dry_run => CommitLogStage::new( - make_kafka_producer(tc), - Topic::new(&tc.physical_topic_name), - Topic::new(source_topic), - consumer_group.to_string(), - ), - _ => CommitLogStage::new( - DryRunProducer, - Topic::new("dry-run-commit-log"), - Topic::new(source_topic), - consumer_group.to_string(), - ), - } + let dest_name = match topic_config { + Some(tc) if !dry_run => tc.physical_topic_name.as_str(), + _ => "dry-run-commit-log", + }; + CommitLogStage::new( + Arc::clone(producer), + Topic::new(dest_name), + Topic::new(source_topic), + consumer_group.to_string(), + ) } -fn make_cogs( +fn make_cogs_stage( + producer: &Arc>, topic_config: &TopicConfig, resource_id: String, dry_run: bool, record_cogs: bool, ) -> CogsStage { - if !dry_run && record_cogs { - CogsStage::new( - make_kafka_producer(topic_config), - Topic::new(&topic_config.physical_topic_name), - resource_id, - ) + let dest_name = if !dry_run && record_cogs { + topic_config.physical_topic_name.as_str() } else { - CogsStage::new(DryRunProducer, Topic::new("dry-run-cogs"), resource_id) - } + "dry-run-cogs" + }; + CogsStage::new(Arc::clone(producer), Topic::new(dest_name), resource_id) +} + +fn make_writer_stage(writer: &Arc) -> ClickHouseWriterStage { + ClickHouseWriterStage::new(Arc::clone(writer)) } fn make_processor_stage( @@ -172,7 +252,7 @@ fn make_processor_stage( ProcessorStage::new(processor, processor_config.clone()) } -fn make_batch( +fn make_batch_stage( max_batch_size: u64, calculation: BatchSizeCalculation, ) -> BatchStage { @@ -183,7 +263,6 @@ fn make_batch( BatchStage::new(PipelineBatchBuffer::new(), max_rows, max_bytes) } -/// Returns (cadence, idle_timeout) for apply_with_timer. fn make_flush_timers( consumer_config: &config::ConsumerConfig, ) -> (Option, Option) { @@ -195,6 +274,7 @@ fn make_flush_timers( // ── Pipeline assembly ────────────────────────────────────── fn make_eap_pipeline( + shared: &SharedResources, consumer_config: &config::ConsumerConfig, storage: &config::StorageConfig, processor_config: &ProcessorConfig, @@ -202,16 +282,12 @@ fn make_eap_pipeline( processing_concurrency: usize, clickhouse_concurrency: usize, dry_run: bool, - dry_run_latency: Option, ) -> EapPipeline { let processor_name = &storage.message_processor.python_class_name; let source_topic_name = &consumer_config.raw_topic.physical_topic_name; assert_eq!(processor_name, "EAPItemsProcessor"); - let insert_columns = Some(crate::processors::eap_items::EAPItemRow::column_names( - processor_config.eap_items_emit_received_at, - )); let resource_id = get_cogs_label(processor_name).unwrap_or_else(|| format!("{}_processor", storage.name)); let (cadence, idle_timeout) = make_flush_timers(consumer_config); @@ -222,27 +298,28 @@ fn make_eap_pipeline( processor_config, ), processing_concurrency, - make_dlq(consumer_config.dlq_topic.as_ref(), dry_run), - make_batch( + make_dlq_handler( + &shared.dlq_producer, + consumer_config.dlq_topic.as_ref(), + dry_run, + ), + make_batch_stage( consumer_config.max_batch_size as u64, consumer_config.max_batch_size_calculation, ), cadence, idle_timeout, - make_writer( - storage, - InsertFormat::RowBinary, - insert_columns, - dry_run_latency, - ), + make_writer_stage(&shared.ch_writer), clickhouse_concurrency, - make_commit_log( + make_commit_log_stage( + &shared.commit_log_producer, consumer_config.commit_log_topic.as_ref(), source_topic_name, consumer_group, dry_run, ), - make_cogs( + make_cogs_stage( + &shared.cogs_producer, &consumer_config.accountant_topic, resource_id, dry_run, @@ -252,26 +329,25 @@ fn make_eap_pipeline( } fn make_faf_pipeline( + shared: &SharedResources, processor: crate::processors::ProcessingFunction, consumer_config: &config::ConsumerConfig, - storage: &config::StorageConfig, processor_config: &ProcessorConfig, processing_concurrency: usize, clickhouse_concurrency: usize, - dry_run_latency: Option, ) -> FireAndForgetPipeline { let (cadence, idle_timeout) = make_flush_timers(consumer_config); FireAndForgetPipeline::new( make_processor_stage(processor, processor_config), processing_concurrency, - make_batch( + make_batch_stage( consumer_config.max_batch_size as u64, consumer_config.max_batch_size_calculation, ), cadence, idle_timeout, - make_writer(storage, InsertFormat::JsonEachRow, None, dry_run_latency), + make_writer_stage(&shared.ch_writer), clickhouse_concurrency, ) } @@ -285,12 +361,17 @@ async fn run_with_rebalance>( loop { let pipeline = build_pipeline(); let mut tracker = OffsetTracker::new(Duration::from_secs(1), source.committer()); - match pipeline.stream(source.stream()).commit(&mut tracker).await { + let result = pipeline.stream(source.stream()).commit(&mut tracker); + + match result.await { Ok(PipelineExit::Rebalance) => { - tracing::info!("Rebalance detected, restarting pipeline"); + tracing::info!("Rebalance detected, restarting pipeline..."); continue; } - Ok(PipelineExit::Shutdown | PipelineExit::Complete) => return 0, + Ok(PipelineExit::Shutdown | PipelineExit::Complete) => { + tracing::info!("Pipeline shutdown"); + return 0; + } Err(e) => { tracing::error!("Pipeline failed: {}", e); return 1; @@ -313,6 +394,7 @@ pub fn pull_consumer( clickhouse_concurrency: usize, max_poll_interval_ms: usize, dry_run_latency_ms: u64, + use_row_binary: bool, ) -> usize { py.allow_threads(|| { pull_consumer_impl( @@ -324,6 +406,7 @@ pub fn pull_consumer( clickhouse_concurrency, max_poll_interval_ms, dry_run_latency_ms, + use_row_binary, ) }) } @@ -338,6 +421,7 @@ fn pull_consumer_impl( clickhouse_concurrency: usize, max_poll_interval_ms: usize, dry_run_latency_ms: u64, + use_row_binary: bool, ) -> usize { setup_logging(); crate::init_sentry_options().expect("failed to initialize sentry-options"); @@ -406,6 +490,7 @@ fn pull_consumer_impl( pipeline = if is_eap { "eap" } else { "fire_and_forget" }, dry_run, dry_run_latency_ms, + use_row_binary, "Starting pull consumer", ); @@ -424,9 +509,19 @@ fn pull_consumer_impl( ); let processor_config = make_processor_config(&storage, &env_config); + let shared = make_shared_resources( + &consumer_config, + &storage, + &processor_config, + use_row_binary, + dry_run, + dry_run_latency, + ); + let result = if is_eap { run_with_rebalance(&source, || { make_eap_pipeline( + &shared, &consumer_config, &storage, &processor_config, @@ -434,20 +529,18 @@ fn pull_consumer_impl( processing_concurrency, clickhouse_concurrency, dry_run, - dry_run_latency, ) }) .await } else { run_with_rebalance(&source, || { make_faf_pipeline( + &shared, processor, &consumer_config, - &storage, &processor_config, processing_concurrency, clickhouse_concurrency, - dry_run_latency, ) }) .await diff --git a/snuba/cli/rust_consumer.py b/snuba/cli/rust_consumer.py index 5f38ce6566..ec89dfe957 100644 --- a/snuba/cli/rust_consumer.py +++ b/snuba/cli/rust_consumer.py @@ -320,6 +320,7 @@ def rust_consumer( clickhouse_concurrency or 2, max_poll_interval_ms, dry_run or 0, + use_row_binary, ) sys.exit(exitcode) From 22aebe12f7f6618e81770310e05bed6f660b63af Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Fri, 28 Aug 2026 16:21:37 -0700 Subject: [PATCH 3/7] ref(pull): implement owned streams. --- rust_snuba/src/pull/pipelines/eap.rs | 20 +++++++++---------- .../src/pull/pipelines/fire_and_forget.rs | 16 +++++++-------- 2 files changed, 18 insertions(+), 18 deletions(-) diff --git a/rust_snuba/src/pull/pipelines/eap.rs b/rust_snuba/src/pull/pipelines/eap.rs index e8fca2b742..753963a079 100644 --- a/rust_snuba/src/pull/pipelines/eap.rs +++ b/rust_snuba/src/pull/pipelines/eap.rs @@ -64,16 +64,16 @@ impl EapPipeline { impl Pipeline for EapPipeline { type Output = BatchMetadata; - fn stream<'a>( - &'a self, - source: impl Stream> + 'a, - ) -> impl Stream> + 'a { + fn stream( + self, + source: impl Stream> + Send, + ) -> impl Stream> + Send { source - .apply_concurrent(&self.processor, self.processing_concurrency) - .apply_with_timer(&self.batch, self.idle_timeout, self.max_batch_time) - .apply_concurrent(&self.writer, self.writer_concurrency) - .apply(&self.commit_log) - .apply(&self.cogs) - .on_reject(&self.dlq_handler) + .apply_concurrent(self.processor, self.processing_concurrency) + .apply_with_timer(self.batch, self.idle_timeout, self.max_batch_time) + .apply_concurrent(self.writer, self.writer_concurrency) + .apply(self.commit_log) + .apply(self.cogs) + .on_reject(self.dlq_handler) } } diff --git a/rust_snuba/src/pull/pipelines/fire_and_forget.rs b/rust_snuba/src/pull/pipelines/fire_and_forget.rs index f6e1a5c132..7cab2e7366 100644 --- a/rust_snuba/src/pull/pipelines/fire_and_forget.rs +++ b/rust_snuba/src/pull/pipelines/fire_and_forget.rs @@ -56,14 +56,14 @@ impl FireAndForgetPipeline { impl Pipeline for FireAndForgetPipeline { type Output = BatchMetadata; - fn stream<'a>( - &'a self, - source: impl Stream> + 'a, - ) -> impl Stream> + 'a { + fn stream( + self, + source: impl Stream> + Send, + ) -> impl Stream> + Send { source - .apply_concurrent(&self.processor, self.processing_concurrency) - .apply_with_timer(&self.batch, self.idle_timeout, self.max_batch_time) - .apply_concurrent(&self.writer, self.writer_concurrency) - .on_reject(&self.rejection_handler) + .apply_concurrent(self.processor, self.processing_concurrency) + .apply_with_timer(self.batch, self.idle_timeout, self.max_batch_time) + .apply_concurrent(self.writer, self.writer_concurrency) + .on_reject(self.rejection_handler) } } From 585dfffbaa70f70b926d214de5460acce177309c Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Fri, 28 Aug 2026 17:03:08 -0700 Subject: [PATCH 4/7] ref(pull): use new pipeline runner. --- rust_snuba/src/pull_consumer.rs | 45 ++++++++------------------------- 1 file changed, 11 insertions(+), 34 deletions(-) diff --git a/rust_snuba/src/pull_consumer.rs b/rust_snuba/src/pull_consumer.rs index dfa6254244..18208f19cb 100644 --- a/rust_snuba/src/pull_consumer.rs +++ b/rust_snuba/src/pull_consumer.rs @@ -8,8 +8,7 @@ use sentry_arroyo::backends::kafka::types::KafkaPayload; use sentry_arroyo::backends::Producer; use sentry_arroyo::metrics; use sentry_arroyo::processing::stream::{ - BatchStage, DlqHandler, KafkaSource, OffsetTracker, Pipeline, PipelineExit, PipelineExt, - PullSource, + BatchStage, DlqHandler, KafkaSource, PipelineRunner, PullSource, }; use sentry_arroyo::types::{Topic, TopicOrPartition}; @@ -18,7 +17,6 @@ use crate::logging::{setup_logging, setup_sentry}; use crate::metrics::statsd::create_dogstatsd_backend; use crate::processors::eap_items::EAPItemRow; use crate::processors::{get_cogs_label, get_processing_function, ProcessingFunctionType}; -use crate::pull::batch::batch_metadata::BatchMetadata; use crate::pull::batch::buffer::PipelineBatchBuffer; use crate::pull::batch::pipeline_batch::PipelineBatch; use crate::pull::pipelines::eap::EapPipeline; @@ -352,34 +350,6 @@ fn make_faf_pipeline( ) } -// ── Rebalance loop ───────────────────────────────────────── - -async fn run_with_rebalance>( - source: &KafkaSource, - build_pipeline: impl Fn() -> P, -) -> usize { - loop { - let pipeline = build_pipeline(); - let mut tracker = OffsetTracker::new(Duration::from_secs(1), source.committer()); - let result = pipeline.stream(source.stream()).commit(&mut tracker); - - match result.await { - Ok(PipelineExit::Rebalance) => { - tracing::info!("Rebalance detected, restarting pipeline..."); - continue; - } - Ok(PipelineExit::Shutdown | PipelineExit::Complete) => { - tracing::info!("Pipeline shutdown"); - return 0; - } - Err(e) => { - tracing::error!("Pipeline failed: {}", e); - return 1; - } - } - } -} - // ── Entry point ──────────────────────────────────────────── #[pyfunction] @@ -519,7 +489,7 @@ fn pull_consumer_impl( ); let result = if is_eap { - run_with_rebalance(&source, || { + PipelineRunner::run(&source, Duration::from_secs(1), || { make_eap_pipeline( &shared, &consumer_config, @@ -533,7 +503,7 @@ fn pull_consumer_impl( }) .await } else { - run_with_rebalance(&source, || { + PipelineRunner::run(&source, Duration::from_secs(1), || { make_faf_pipeline( &shared, processor, @@ -547,7 +517,14 @@ fn pull_consumer_impl( }; source.shutdown(); - result + + match result { + Ok(()) => 0, + Err(e) => { + tracing::error!("Pipeline failed: {e}"); + 1 + } + } }); exit_code From 4ef6dd3252ceb2354069de847d1e2238f1492b99 Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Tue, 1 Sep 2026 10:25:37 -0700 Subject: [PATCH 5/7] ref(pull): Re-wire to branch arroyo. --- rust_snuba/Cargo.lock | 3 ++- rust_snuba/Cargo.toml | 5 +---- 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/rust_snuba/Cargo.lock b/rust_snuba/Cargo.lock index 34386b24c4..c907096b10 100644 --- a/rust_snuba/Cargo.lock +++ b/rust_snuba/Cargo.lock @@ -3911,7 +3911,8 @@ dependencies = [ [[package]] name = "sentry_arroyo" -version = "2.43.0" +version = "2.43.1" +source = "git+https://github.com/getsentry/arroyo?rev=1bf5057#1bf50571d5492cf6c208341aa8a95342cf88177b" dependencies = [ "async-stream", "chrono", diff --git a/rust_snuba/Cargo.toml b/rust_snuba/Cargo.toml index 503ebae250..bc714c8589 100644 --- a/rust_snuba/Cargo.toml +++ b/rust_snuba/Cargo.toml @@ -69,7 +69,7 @@ sentry = { version = "0.48.5", default-features = false, features = [ sentry-kafka-schemas = "2.1.39" sentry-options = "1.2.1" sentry_protos = "0.60.0" -sentry_arroyo = { git = "https://github.com/getsentry/arroyo", branch = "rbroughan/pull-based-arroyo-part-2", features = ["ssl"] } +sentry_arroyo = { git = "https://github.com/getsentry/arroyo", rev = "1bf5057", features = ["ssl"] } sentry_usage_accountant = { version = "0.1.2", features = ["kafka"] } seq-macro = "0.3" serde = { version = "1.0", features = ["derive"] } @@ -104,6 +104,3 @@ tikv-jemallocator = "0.5" [[bench]] name = "processors" harness = false - -[patch."https://github.com/getsentry/arroyo"] -sentry_arroyo = { path = "/Users/ryanbroughan/code/arroyo" } From 185230f63294d1e7f6bae171fffd727d3ad1130f Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Fri, 11 Sep 2026 13:44:13 -0700 Subject: [PATCH 6/7] ref(pull): break up factories and fix some arroyo callsites. --- rust_snuba/Cargo.lock | 2 +- rust_snuba/Cargo.toml | 2 +- rust_snuba/src/accepted_outcomes_consumer.rs | 6 +- rust_snuba/src/consumer.rs | 6 +- rust_snuba/src/factory_v2.rs | 3 +- rust_snuba/src/pull_consumer.rs | 531 ------------------- rust_snuba/src/pull_consumer/entrypoint.rs | 191 +++++++ rust_snuba/src/pull_consumer/factories.rs | 249 +++++++++ rust_snuba/src/pull_consumer/mod.rs | 19 + rust_snuba/src/pull_consumer/pipelines.rs | 86 +++ 10 files changed, 557 insertions(+), 538 deletions(-) delete mode 100644 rust_snuba/src/pull_consumer.rs create mode 100644 rust_snuba/src/pull_consumer/entrypoint.rs create mode 100644 rust_snuba/src/pull_consumer/factories.rs create mode 100644 rust_snuba/src/pull_consumer/mod.rs create mode 100644 rust_snuba/src/pull_consumer/pipelines.rs diff --git a/rust_snuba/Cargo.lock b/rust_snuba/Cargo.lock index c907096b10..0719495a01 100644 --- a/rust_snuba/Cargo.lock +++ b/rust_snuba/Cargo.lock @@ -3912,7 +3912,7 @@ dependencies = [ [[package]] name = "sentry_arroyo" version = "2.43.1" -source = "git+https://github.com/getsentry/arroyo?rev=1bf5057#1bf50571d5492cf6c208341aa8a95342cf88177b" +source = "git+https://github.com/getsentry/arroyo?branch=rbroughan%2Fpull-based-arroyo-part-2#1bf50571d5492cf6c208341aa8a95342cf88177b" dependencies = [ "async-stream", "chrono", diff --git a/rust_snuba/Cargo.toml b/rust_snuba/Cargo.toml index bc714c8589..da627a2a9b 100644 --- a/rust_snuba/Cargo.toml +++ b/rust_snuba/Cargo.toml @@ -69,7 +69,7 @@ sentry = { version = "0.48.5", default-features = false, features = [ sentry-kafka-schemas = "2.1.39" sentry-options = "1.2.1" sentry_protos = "0.60.0" -sentry_arroyo = { git = "https://github.com/getsentry/arroyo", rev = "1bf5057", features = ["ssl"] } +sentry_arroyo = { git = "https://github.com/getsentry/arroyo", branch = "rbroughan/pull-based-arroyo-part-2", features = ["ssl"] } sentry_usage_accountant = { version = "0.1.2", features = ["kafka"] } seq-macro = "0.3" serde = { version = "1.0", features = ["derive"] } diff --git a/rust_snuba/src/accepted_outcomes_consumer.rs b/rust_snuba/src/accepted_outcomes_consumer.rs index e3e29988de..4f8b2f37da 100644 --- a/rust_snuba/src/accepted_outcomes_consumer.rs +++ b/rust_snuba/src/accepted_outcomes_consumer.rs @@ -180,7 +180,8 @@ pub fn accepted_outcomes_consumer_impl( let dlq_policy = consumer_config.dlq_topic.map(|dlq_topic_config| { let dlq_producer_config = KafkaConfig::new_producer_config(vec![], Some(dlq_topic_config.broker_config)); - let dlq_producer = KafkaProducer::new(dlq_producer_config); + let dlq_producer = + KafkaProducer::new(dlq_producer_config).expect("failed to create kafka producer"); let kafka_dlq_producer = Box::new(KafkaDlqProducer::new( dlq_producer, @@ -219,7 +220,8 @@ pub fn accepted_outcomes_consumer_impl( // TODO consider adding a higher linger.ms to the config for batching produced messages let producer_config = KafkaConfig::new_producer_config(vec![], Some(topic_config.broker_config)); - let producer = Arc::new(KafkaProducer::new(producer_config)); + let producer = + Arc::new(KafkaProducer::new(producer_config).expect("failed to create kafka producer")); let produce_topic = Topic::new(&topic_config.physical_topic_name); let factory = AcceptedOutcomesStrategyFactory { diff --git a/rust_snuba/src/consumer.rs b/rust_snuba/src/consumer.rs index 9945ed9ab1..d42555a2ec 100644 --- a/rust_snuba/src/consumer.rs +++ b/rust_snuba/src/consumer.rs @@ -224,7 +224,8 @@ pub fn consumer_impl( let producer = KafkaProducer::new(KafkaConfig::new_producer_config( vec![], Some(dlq_topic_config.broker_config), - )); + )) + .expect("failed to create kafka producer"); let kafka_dlq_producer = Box::new(KafkaDlqProducer::new( producer, @@ -252,7 +253,8 @@ pub fn consumer_impl( } else if let Some(topic_config) = consumer_config.commit_log_topic { let producer_config = KafkaConfig::new_producer_config(vec![], Some(topic_config.broker_config)); - let producer = KafkaProducer::new(producer_config); + let producer = + KafkaProducer::new(producer_config).expect("failed to create kafka producer"); Some(( Arc::new(producer), Topic::new(&topic_config.physical_topic_name), diff --git a/rust_snuba/src/factory_v2.rs b/rust_snuba/src/factory_v2.rs index 1d3f006ec6..5ec78bf6fe 100644 --- a/rust_snuba/src/factory_v2.rs +++ b/rust_snuba/src/factory_v2.rs @@ -246,7 +246,8 @@ impl ProcessingStrategyFactory for ConsumerStrategyFactoryV2 { self.replacements_config.clone().expect( "replacements topic required for processors that emit replacements", ); - let producer = KafkaProducer::new(replacements_config); + let producer = KafkaProducer::new(replacements_config) + .expect("failed to create kafka producer"); ProduceReplacements::new( next_step, producer, diff --git a/rust_snuba/src/pull_consumer.rs b/rust_snuba/src/pull_consumer.rs deleted file mode 100644 index 18208f19cb..0000000000 --- a/rust_snuba/src/pull_consumer.rs +++ /dev/null @@ -1,531 +0,0 @@ -use std::sync::Arc; -use std::time::Duration; - -use pyo3::prelude::*; -use sentry_arroyo::backends::kafka::config::KafkaConfig; -use sentry_arroyo::backends::kafka::producer::KafkaProducer; -use sentry_arroyo::backends::kafka::types::KafkaPayload; -use sentry_arroyo::backends::Producer; -use sentry_arroyo::metrics; -use sentry_arroyo::processing::stream::{ - BatchStage, DlqHandler, KafkaSource, PipelineRunner, PullSource, -}; -use sentry_arroyo::types::{Topic, TopicOrPartition}; - -use crate::config::{self, BatchSizeCalculation, ProcessorConfig, TopicConfig}; -use crate::logging::{setup_logging, setup_sentry}; -use crate::metrics::statsd::create_dogstatsd_backend; -use crate::processors::eap_items::EAPItemRow; -use crate::processors::{get_cogs_label, get_processing_function, ProcessingFunctionType}; -use crate::pull::batch::buffer::PipelineBatchBuffer; -use crate::pull::batch::pipeline_batch::PipelineBatch; -use crate::pull::pipelines::eap::EapPipeline; -use crate::pull::pipelines::fire_and_forget::FireAndForgetPipeline; -use crate::pull::producers::DryRunProducer; -use crate::pull::stages::clickhouse_writer_stage::ClickHouseWriterStage; -use crate::pull::stages::cogs_stage::CogsStage; -use crate::pull::stages::commit_log_stage::CommitLogStage; -use crate::pull::stages::processor_stage::ProcessorStage; -use crate::pull::writer::{ClickHouseWriter, DryRunWriter}; -use crate::strategies::clickhouse::writer_v2::{ClickhouseClient, InsertFormat}; - -/// Allowed processors for the fire-and-forget pipeline. -const FIRE_AND_FORGET_PROCESSORS: &[&str] = &[ - "FunctionsMessageProcessor", - "ProfilesMessageProcessor", - "QuerylogProcessor", - "ReplaysProcessor", - "OutcomesProcessor", - "ProfileChunksProcessor", - "LlmProxyCostProcessor", -]; - -/// Allowed processors for the EAP pipeline. -const EAP_PROCESSORS: &[&str] = &["EAPItemsProcessor"]; - -// ── Shared resources ────────────────────────────────────── - -struct SharedResources { - dlq_producer: Arc>, - commit_log_producer: Arc>, - cogs_producer: Arc>, - ch_writer: Arc, -} - -fn make_shared_resources( - consumer_config: &config::ConsumerConfig, - storage: &config::StorageConfig, - processor_config: &ProcessorConfig, - use_row_binary: bool, - dry_run: bool, - dry_run_latency: Option, -) -> SharedResources { - SharedResources { - dlq_producer: Arc::from(make_dlq_producer(consumer_config, dry_run)), - commit_log_producer: Arc::from(make_commit_log_producer(consumer_config, dry_run)), - cogs_producer: Arc::from(make_cogs_producer(consumer_config, dry_run)), - ch_writer: Arc::from(make_ch_writer( - storage, - processor_config, - use_row_binary, - dry_run_latency, - )), - } -} - -// ── Producer factories ──────────────────── - -fn make_kafka_producer(topic_config: &TopicConfig) -> KafkaProducer { - KafkaProducer::new(KafkaConfig::new_producer_config( - vec![], - Some(topic_config.broker_config.clone()), - )) -} - -fn make_dlq_producer( - consumer_config: &config::ConsumerConfig, - dry_run: bool, -) -> Box> { - match consumer_config.dlq_topic.as_ref() { - Some(tc) if !dry_run => Box::new(make_kafka_producer(tc)), - _ => Box::new(DryRunProducer), - } -} - -fn make_commit_log_producer( - consumer_config: &config::ConsumerConfig, - dry_run: bool, -) -> Box> { - match consumer_config.commit_log_topic.as_ref() { - Some(tc) if !dry_run => Box::new(make_kafka_producer(tc)), - _ => Box::new(DryRunProducer), - } -} - -fn make_cogs_producer( - consumer_config: &config::ConsumerConfig, - dry_run: bool, -) -> Box> { - if !dry_run && consumer_config.env.record_cogs { - Box::new(make_kafka_producer(&consumer_config.accountant_topic)) - } else { - Box::new(DryRunProducer) - } -} - -fn make_ch_writer( - storage: &config::StorageConfig, - processor_config: &ProcessorConfig, - use_row_binary: bool, - dry_run_latency: Option, -) -> Box { - if let Some(latency) = dry_run_latency { - return Box::new(DryRunWriter::new(latency)); - } - - if use_row_binary { - Box::new(ClickhouseClient::new( - &storage.clickhouse_cluster, - &storage.clickhouse_table_name, - storage.name.clone(), - InsertFormat::RowBinary, - Some(EAPItemRow::column_names( - processor_config.eap_items_emit_received_at, - )), - )) - } else { - Box::new(ClickhouseClient::new( - &storage.clickhouse_cluster, - &storage.clickhouse_table_name, - storage.name.clone(), - InsertFormat::JsonEachRow, - None, - )) - } -} - -// ── Stage factories ─────────────────────────────────────── - -fn resolve_processor( - processor_name: &str, -) -> Result { - match get_processing_function(processor_name) { - Some(ProcessingFunctionType::ProcessingFunction(f)) => Ok(f), - Some(ProcessingFunctionType::ProcessingFunctionWithReplacements(_)) => Err(format!( - "{processor_name} is a replacement processor — not supported" - )), - None => Err(format!("Unknown processor: {processor_name}")), - } -} - -fn make_kafka_source( - consumer_config: &config::ConsumerConfig, - consumer_group: &str, - auto_offset_reset: &str, - no_strict_offset_reset: bool, - max_poll_interval_ms: usize, -) -> KafkaSource { - let kafka_config = KafkaConfig::new_consumer_config( - vec![], - consumer_group.to_owned(), - auto_offset_reset - .parse() - .expect("Invalid auto_offset_reset"), - !no_strict_offset_reset, - max_poll_interval_ms, - Some(consumer_config.raw_topic.broker_config.clone()), - ); - let topic = Topic::new(&consumer_config.raw_topic.physical_topic_name); - KafkaSource::new(kafka_config, &[topic]) -} - -fn make_processor_config( - storage: &config::StorageConfig, - env_config: &config::EnvConfig, -) -> ProcessorConfig { - ProcessorConfig { - env_config: env_config.clone(), - storage_name: storage.name.clone(), - eap_items_emit_received_at: crate::processors::eap_items::emit_received_at(), - } -} - -fn make_dlq_handler( - producer: &Arc>, - topic_config: Option<&TopicConfig>, - dry_run: bool, -) -> DlqHandler { - let topic_name = match topic_config { - Some(tc) if !dry_run => tc.physical_topic_name.as_str(), - _ => "dry-run-dlq", - }; - DlqHandler::new( - Arc::clone(producer), - TopicOrPartition::Topic(Topic::new(topic_name)), - ) -} - -fn make_commit_log_stage( - producer: &Arc>, - topic_config: Option<&TopicConfig>, - source_topic: &str, - consumer_group: &str, - dry_run: bool, -) -> CommitLogStage { - let dest_name = match topic_config { - Some(tc) if !dry_run => tc.physical_topic_name.as_str(), - _ => "dry-run-commit-log", - }; - CommitLogStage::new( - Arc::clone(producer), - Topic::new(dest_name), - Topic::new(source_topic), - consumer_group.to_string(), - ) -} - -fn make_cogs_stage( - producer: &Arc>, - topic_config: &TopicConfig, - resource_id: String, - dry_run: bool, - record_cogs: bool, -) -> CogsStage { - let dest_name = if !dry_run && record_cogs { - topic_config.physical_topic_name.as_str() - } else { - "dry-run-cogs" - }; - CogsStage::new(Arc::clone(producer), Topic::new(dest_name), resource_id) -} - -fn make_writer_stage(writer: &Arc) -> ClickHouseWriterStage { - ClickHouseWriterStage::new(Arc::clone(writer)) -} - -fn make_processor_stage( - processor: crate::processors::ProcessingFunction, - processor_config: &ProcessorConfig, -) -> ProcessorStage { - ProcessorStage::new(processor, processor_config.clone()) -} - -fn make_batch_stage( - max_batch_size: u64, - calculation: BatchSizeCalculation, -) -> BatchStage { - let (max_rows, max_bytes) = match calculation { - BatchSizeCalculation::Rows => (max_batch_size, u64::MAX), - BatchSizeCalculation::Bytes => (u64::MAX, max_batch_size), - }; - BatchStage::new(PipelineBatchBuffer::new(), max_rows, max_bytes) -} - -fn make_flush_timers( - consumer_config: &config::ConsumerConfig, -) -> (Option, Option) { - let cadence = Some(Duration::from_millis(consumer_config.max_batch_time_ms)); - let idle_timeout = None; - (cadence, idle_timeout) -} - -// ── Pipeline assembly ────────────────────────────────────── - -fn make_eap_pipeline( - shared: &SharedResources, - consumer_config: &config::ConsumerConfig, - storage: &config::StorageConfig, - processor_config: &ProcessorConfig, - consumer_group: &str, - processing_concurrency: usize, - clickhouse_concurrency: usize, - dry_run: bool, -) -> EapPipeline { - let processor_name = &storage.message_processor.python_class_name; - let source_topic_name = &consumer_config.raw_topic.physical_topic_name; - - assert_eq!(processor_name, "EAPItemsProcessor"); - - let resource_id = - get_cogs_label(processor_name).unwrap_or_else(|| format!("{}_processor", storage.name)); - let (cadence, idle_timeout) = make_flush_timers(consumer_config); - - EapPipeline::new( - make_processor_stage( - crate::processors::eap_items::process_message_row_binary, - processor_config, - ), - processing_concurrency, - make_dlq_handler( - &shared.dlq_producer, - consumer_config.dlq_topic.as_ref(), - dry_run, - ), - make_batch_stage( - consumer_config.max_batch_size as u64, - consumer_config.max_batch_size_calculation, - ), - cadence, - idle_timeout, - make_writer_stage(&shared.ch_writer), - clickhouse_concurrency, - make_commit_log_stage( - &shared.commit_log_producer, - consumer_config.commit_log_topic.as_ref(), - source_topic_name, - consumer_group, - dry_run, - ), - make_cogs_stage( - &shared.cogs_producer, - &consumer_config.accountant_topic, - resource_id, - dry_run, - consumer_config.env.record_cogs, - ), - ) -} - -fn make_faf_pipeline( - shared: &SharedResources, - processor: crate::processors::ProcessingFunction, - consumer_config: &config::ConsumerConfig, - processor_config: &ProcessorConfig, - processing_concurrency: usize, - clickhouse_concurrency: usize, -) -> FireAndForgetPipeline { - let (cadence, idle_timeout) = make_flush_timers(consumer_config); - - FireAndForgetPipeline::new( - make_processor_stage(processor, processor_config), - processing_concurrency, - make_batch_stage( - consumer_config.max_batch_size as u64, - consumer_config.max_batch_size_calculation, - ), - cadence, - idle_timeout, - make_writer_stage(&shared.ch_writer), - clickhouse_concurrency, - ) -} - -// ── Entry point ──────────────────────────────────────────── - -#[pyfunction] -#[allow(clippy::too_many_arguments)] -pub fn pull_consumer( - py: Python<'_>, - consumer_group: &str, - auto_offset_reset: &str, - no_strict_offset_reset: bool, - consumer_config_raw: &str, - processing_concurrency: usize, - clickhouse_concurrency: usize, - max_poll_interval_ms: usize, - dry_run_latency_ms: u64, - use_row_binary: bool, -) -> usize { - py.allow_threads(|| { - pull_consumer_impl( - consumer_group, - auto_offset_reset, - no_strict_offset_reset, - consumer_config_raw, - processing_concurrency, - clickhouse_concurrency, - max_poll_interval_ms, - dry_run_latency_ms, - use_row_binary, - ) - }) -} - -#[allow(clippy::too_many_arguments)] -fn pull_consumer_impl( - consumer_group: &str, - auto_offset_reset: &str, - no_strict_offset_reset: bool, - consumer_config_raw: &str, - processing_concurrency: usize, - clickhouse_concurrency: usize, - max_poll_interval_ms: usize, - dry_run_latency_ms: u64, - use_row_binary: bool, -) -> usize { - setup_logging(); - crate::init_sentry_options().expect("failed to initialize sentry-options"); - - let consumer_config = config::ConsumerConfig::load_from_str(consumer_config_raw) - .expect("failed to parse consumer config"); - - assert_eq!( - consumer_config.storages.len(), - 1, - "pull consumer only supports a single storage" - ); - - let storage = consumer_config.storages[0].clone(); - let processor_name = storage.message_processor.python_class_name.clone(); - let env_config = consumer_config.env.clone(); - - // Sentry - let mut _sentry_guard = None; - if let Some(ref dsn) = consumer_config.env.sentry_dsn { - std::env::set_var("RUST_BACKTRACE", "1"); - _sentry_guard = Some(setup_sentry(dsn)); - } - - // Metrics - { - let tags = [ - ("storage", storage.name.clone()), - ("consumer_group", consumer_group.to_owned()), - ]; - sentry::configure_scope(|scope| { - scope.set_tag("storage", &storage.name); - scope.set_tag("consumer_group", consumer_group); - }); - if let Some(backend) = create_dogstatsd_backend(&env_config, "snuba.consumer", &tags) { - metrics::init(backend).unwrap(); - } - } - - let processor = match resolve_processor(&processor_name) { - Ok(f) => f, - Err(msg) => { - tracing::error!("{msg}"); - return 1; - } - }; - - // Validate pipeline type - let is_eap = EAP_PROCESSORS.contains(&processor_name.as_str()); - let is_faf = FIRE_AND_FORGET_PROCESSORS.contains(&processor_name.as_str()); - if !is_eap && !is_faf { - tracing::error!("{processor_name} is not supported by the pull consumer"); - return 1; - } - - let dry_run = dry_run_latency_ms > 0; - let dry_run_latency = if dry_run { - Some(Duration::from_millis(dry_run_latency_ms)) - } else { - None - }; - - tracing::info!( - storage = storage.name, - processor = processor_name.as_str(), - pipeline = if is_eap { "eap" } else { "fire_and_forget" }, - dry_run, - dry_run_latency_ms, - use_row_binary, - "Starting pull consumer", - ); - - let rt = tokio::runtime::Builder::new_multi_thread() - .enable_all() - .build() - .expect("failed to build tokio runtime"); - - let exit_code = rt.block_on(async { - let source = make_kafka_source( - &consumer_config, - consumer_group, - auto_offset_reset, - no_strict_offset_reset, - max_poll_interval_ms, - ); - let processor_config = make_processor_config(&storage, &env_config); - - let shared = make_shared_resources( - &consumer_config, - &storage, - &processor_config, - use_row_binary, - dry_run, - dry_run_latency, - ); - - let result = if is_eap { - PipelineRunner::run(&source, Duration::from_secs(1), || { - make_eap_pipeline( - &shared, - &consumer_config, - &storage, - &processor_config, - consumer_group, - processing_concurrency, - clickhouse_concurrency, - dry_run, - ) - }) - .await - } else { - PipelineRunner::run(&source, Duration::from_secs(1), || { - make_faf_pipeline( - &shared, - processor, - &consumer_config, - &processor_config, - processing_concurrency, - clickhouse_concurrency, - ) - }) - .await - }; - - source.shutdown(); - - match result { - Ok(()) => 0, - Err(e) => { - tracing::error!("Pipeline failed: {e}"); - 1 - } - } - }); - - exit_code -} diff --git a/rust_snuba/src/pull_consumer/entrypoint.rs b/rust_snuba/src/pull_consumer/entrypoint.rs new file mode 100644 index 0000000000..3f94ce236b --- /dev/null +++ b/rust_snuba/src/pull_consumer/entrypoint.rs @@ -0,0 +1,191 @@ +use std::time::Duration; + +use pyo3::prelude::*; +use sentry_arroyo::metrics; +use sentry_arroyo::processing::stream::{PipelineRunner, PullSource}; + +use crate::config; +use crate::logging::{setup_logging, setup_sentry}; +use crate::metrics::statsd::create_dogstatsd_backend; + +use super::factories; +use super::pipelines; +use super::{EAP_PROCESSORS, FIRE_AND_FORGET_PROCESSORS}; + +#[pyfunction] +#[allow(clippy::too_many_arguments)] +pub fn pull_consumer( + py: Python<'_>, + consumer_group: &str, + auto_offset_reset: &str, + no_strict_offset_reset: bool, + consumer_config_raw: &str, + processing_concurrency: usize, + clickhouse_concurrency: usize, + max_poll_interval_ms: usize, + dry_run_latency_ms: u64, + use_row_binary: bool, +) -> usize { + py.allow_threads(|| { + pull_consumer_impl( + consumer_group, + auto_offset_reset, + no_strict_offset_reset, + consumer_config_raw, + processing_concurrency, + clickhouse_concurrency, + max_poll_interval_ms, + dry_run_latency_ms, + use_row_binary, + ) + }) +} + +#[allow(clippy::too_many_arguments)] +fn pull_consumer_impl( + consumer_group: &str, + auto_offset_reset: &str, + no_strict_offset_reset: bool, + consumer_config_raw: &str, + processing_concurrency: usize, + clickhouse_concurrency: usize, + max_poll_interval_ms: usize, + dry_run_latency_ms: u64, + use_row_binary: bool, +) -> usize { + setup_logging(); + crate::init_sentry_options().expect("failed to initialize sentry-options"); + + let consumer_config = config::ConsumerConfig::load_from_str(consumer_config_raw) + .expect("failed to parse consumer config"); + + assert_eq!( + consumer_config.storages.len(), + 1, + "pull consumer only supports a single storage" + ); + + let storage = consumer_config.storages[0].clone(); + let processor_name = storage.message_processor.python_class_name.clone(); + let env_config = consumer_config.env.clone(); + + // Sentry + let mut _sentry_guard = None; + if let Some(ref dsn) = consumer_config.env.sentry_dsn { + std::env::set_var("RUST_BACKTRACE", "1"); + _sentry_guard = Some(setup_sentry(dsn)); + } + + // Metrics + { + let tags = [ + ("storage", storage.name.clone()), + ("consumer_group", consumer_group.to_owned()), + ]; + sentry::configure_scope(|scope| { + scope.set_tag("storage", &storage.name); + scope.set_tag("consumer_group", consumer_group); + }); + if let Some(backend) = create_dogstatsd_backend(&env_config, "snuba.consumer", &tags) { + metrics::init(backend).unwrap(); + } + } + + let processor = match factories::resolve_processor(&processor_name) { + Ok(f) => f, + Err(msg) => { + tracing::error!("{msg}"); + return 1; + } + }; + + // Validate pipeline type + let is_eap = EAP_PROCESSORS.contains(&processor_name.as_str()); + let is_faf = FIRE_AND_FORGET_PROCESSORS.contains(&processor_name.as_str()); + if !is_eap && !is_faf { + tracing::error!("{processor_name} is not supported by the pull consumer"); + return 1; + } + + let dry_run = dry_run_latency_ms > 0; + let dry_run_latency = if dry_run { + Some(Duration::from_millis(dry_run_latency_ms)) + } else { + None + }; + + tracing::info!( + storage = storage.name, + processor = processor_name.as_str(), + pipeline = if is_eap { "eap" } else { "fire_and_forget" }, + dry_run, + dry_run_latency_ms, + use_row_binary, + "Starting pull consumer", + ); + + let rt = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .build() + .expect("failed to build tokio runtime"); + + let exit_code = rt.block_on(async { + let source = factories::make_kafka_source( + &consumer_config, + consumer_group, + auto_offset_reset, + no_strict_offset_reset, + max_poll_interval_ms, + ); + let processor_config = factories::make_processor_config(&storage, &env_config); + + let shared = factories::make_shared_resources( + &consumer_config, + &storage, + &processor_config, + use_row_binary, + dry_run, + dry_run_latency, + ); + + let result = if is_eap { + PipelineRunner::run(&source, Duration::from_secs(1), || { + pipelines::make_eap_pipeline( + &shared, + &consumer_config, + &storage, + &processor_config, + consumer_group, + processing_concurrency, + clickhouse_concurrency, + dry_run, + ) + }) + .await + } else { + PipelineRunner::run(&source, Duration::from_secs(1), || { + pipelines::make_faf_pipeline( + &shared, + processor, + &consumer_config, + &processor_config, + processing_concurrency, + clickhouse_concurrency, + ) + }) + .await + }; + + source.shutdown(); + + match result { + Ok(()) => 0, + Err(e) => { + tracing::error!("Pipeline failed: {e}"); + 1 + } + } + }); + + exit_code +} diff --git a/rust_snuba/src/pull_consumer/factories.rs b/rust_snuba/src/pull_consumer/factories.rs new file mode 100644 index 0000000000..867c01c34b --- /dev/null +++ b/rust_snuba/src/pull_consumer/factories.rs @@ -0,0 +1,249 @@ +use std::sync::Arc; +use std::time::Duration; + +use sentry_arroyo::backends::kafka::config::KafkaConfig; +use sentry_arroyo::backends::kafka::producer::KafkaProducer; +use sentry_arroyo::backends::kafka::types::KafkaPayload; +use sentry_arroyo::backends::Producer; +use sentry_arroyo::processing::stream::{BatchStage, DlqHandler, KafkaSource}; +use sentry_arroyo::types::{Topic, TopicOrPartition}; + +use crate::config::{self, BatchSizeCalculation, ProcessorConfig, TopicConfig}; +use crate::processors::eap_items::EAPItemRow; +use crate::processors::{get_processing_function, ProcessingFunctionType}; +use crate::pull::batch::buffer::PipelineBatchBuffer; +use crate::pull::batch::pipeline_batch::PipelineBatch; +use crate::pull::producers::DryRunProducer; +use crate::pull::stages::clickhouse_writer_stage::ClickHouseWriterStage; +use crate::pull::stages::cogs_stage::CogsStage; +use crate::pull::stages::commit_log_stage::CommitLogStage; +use crate::pull::stages::processor_stage::ProcessorStage; +use crate::pull::writer::{ClickHouseWriter, DryRunWriter}; +use crate::strategies::clickhouse::writer_v2::{ClickhouseClient, InsertFormat}; + +// ── Shared resources ────────────────────────────────────── + +pub(super) struct SharedResources { + pub dlq_producer: Arc>, + pub commit_log_producer: Arc>, + pub cogs_producer: Arc>, + pub ch_writer: Arc, +} + +pub(super) fn make_shared_resources( + consumer_config: &config::ConsumerConfig, + storage: &config::StorageConfig, + processor_config: &ProcessorConfig, + use_row_binary: bool, + dry_run: bool, + dry_run_latency: Option, +) -> SharedResources { + SharedResources { + dlq_producer: Arc::from(make_dlq_producer(consumer_config, dry_run)), + commit_log_producer: Arc::from(make_commit_log_producer(consumer_config, dry_run)), + cogs_producer: Arc::from(make_cogs_producer(consumer_config, dry_run)), + ch_writer: Arc::from(make_ch_writer( + storage, + processor_config, + use_row_binary, + dry_run_latency, + )), + } +} + +// ── Producer factories ──────────────────── + +fn make_kafka_producer(topic_config: &TopicConfig) -> KafkaProducer { + KafkaProducer::new(KafkaConfig::new_producer_config( + vec![], + Some(topic_config.broker_config.clone()), + )) + .expect("failed to create kafka producer") +} + +fn make_dlq_producer( + consumer_config: &config::ConsumerConfig, + dry_run: bool, +) -> Box> { + match consumer_config.dlq_topic.as_ref() { + Some(tc) if !dry_run => Box::new(make_kafka_producer(tc)), + _ => Box::new(DryRunProducer), + } +} + +fn make_commit_log_producer( + consumer_config: &config::ConsumerConfig, + dry_run: bool, +) -> Box> { + match consumer_config.commit_log_topic.as_ref() { + Some(tc) if !dry_run => Box::new(make_kafka_producer(tc)), + _ => Box::new(DryRunProducer), + } +} + +fn make_cogs_producer( + consumer_config: &config::ConsumerConfig, + dry_run: bool, +) -> Box> { + if !dry_run && consumer_config.env.record_cogs { + Box::new(make_kafka_producer(&consumer_config.accountant_topic)) + } else { + Box::new(DryRunProducer) + } +} + +fn make_ch_writer( + storage: &config::StorageConfig, + processor_config: &ProcessorConfig, + use_row_binary: bool, + dry_run_latency: Option, +) -> Box { + if let Some(latency) = dry_run_latency { + return Box::new(DryRunWriter::new(latency)); + } + + if use_row_binary { + Box::new(ClickhouseClient::new( + &storage.clickhouse_cluster, + &storage.clickhouse_table_name, + storage.name.clone(), + InsertFormat::RowBinary, + Some(EAPItemRow::column_names( + processor_config.eap_items_emit_received_at, + )), + )) + } else { + Box::new(ClickhouseClient::new( + &storage.clickhouse_cluster, + &storage.clickhouse_table_name, + storage.name.clone(), + InsertFormat::JsonEachRow, + None, + )) + } +} + +// ── Stage factories ─────────────────────────────────────── + +pub(super) fn resolve_processor( + processor_name: &str, +) -> Result { + match get_processing_function(processor_name) { + Some(ProcessingFunctionType::ProcessingFunction(f)) => Ok(f), + Some(ProcessingFunctionType::ProcessingFunctionWithReplacements(_)) => Err(format!( + "{processor_name} is a replacement processor — not supported" + )), + None => Err(format!("Unknown processor: {processor_name}")), + } +} + +pub(super) fn make_kafka_source( + consumer_config: &config::ConsumerConfig, + consumer_group: &str, + auto_offset_reset: &str, + no_strict_offset_reset: bool, + max_poll_interval_ms: usize, +) -> KafkaSource { + let kafka_config = KafkaConfig::new_consumer_config( + vec![], + consumer_group.to_owned(), + auto_offset_reset + .parse() + .expect("Invalid auto_offset_reset"), + !no_strict_offset_reset, + max_poll_interval_ms, + Some(consumer_config.raw_topic.broker_config.clone()), + ); + let topic = Topic::new(&consumer_config.raw_topic.physical_topic_name); + KafkaSource::new(kafka_config, &[topic]) +} + +pub(super) fn make_processor_config( + storage: &config::StorageConfig, + env_config: &config::EnvConfig, +) -> ProcessorConfig { + ProcessorConfig { + env_config: env_config.clone(), + storage_name: storage.name.clone(), + eap_items_emit_received_at: crate::processors::eap_items::emit_received_at(), + } +} + +pub(super) fn make_dlq_handler( + producer: &Arc>, + topic_config: Option<&TopicConfig>, + dry_run: bool, +) -> DlqHandler { + let topic_name = match topic_config { + Some(tc) if !dry_run => tc.physical_topic_name.as_str(), + _ => "dry-run-dlq", + }; + DlqHandler::new( + Arc::clone(producer), + TopicOrPartition::Topic(Topic::new(topic_name)), + ) +} + +pub(super) fn make_commit_log_stage( + producer: &Arc>, + topic_config: Option<&TopicConfig>, + source_topic: &str, + consumer_group: &str, + dry_run: bool, +) -> CommitLogStage { + let dest_name = match topic_config { + Some(tc) if !dry_run => tc.physical_topic_name.as_str(), + _ => "dry-run-commit-log", + }; + CommitLogStage::new( + Arc::clone(producer), + Topic::new(dest_name), + Topic::new(source_topic), + consumer_group.to_string(), + ) +} + +pub(super) fn make_cogs_stage( + producer: &Arc>, + topic_config: &TopicConfig, + resource_id: String, + dry_run: bool, + record_cogs: bool, +) -> CogsStage { + let dest_name = if !dry_run && record_cogs { + topic_config.physical_topic_name.as_str() + } else { + "dry-run-cogs" + }; + CogsStage::new(Arc::clone(producer), Topic::new(dest_name), resource_id) +} + +pub(super) fn make_writer_stage(writer: &Arc) -> ClickHouseWriterStage { + ClickHouseWriterStage::new(Arc::clone(writer)) +} + +pub(super) fn make_processor_stage( + processor: crate::processors::ProcessingFunction, + processor_config: &ProcessorConfig, +) -> ProcessorStage { + ProcessorStage::new(processor, processor_config.clone()) +} + +pub(super) fn make_batch_stage( + max_batch_size: u64, + calculation: BatchSizeCalculation, +) -> BatchStage { + let (max_rows, max_bytes) = match calculation { + BatchSizeCalculation::Rows => (max_batch_size, u64::MAX), + BatchSizeCalculation::Bytes => (u64::MAX, max_batch_size), + }; + BatchStage::new(PipelineBatchBuffer::new(), max_rows, max_bytes) +} + +pub(super) fn make_flush_timers( + consumer_config: &config::ConsumerConfig, +) -> (Option, Option) { + let cadence = Some(Duration::from_millis(consumer_config.max_batch_time_ms)); + let idle_timeout = None; + (cadence, idle_timeout) +} diff --git a/rust_snuba/src/pull_consumer/mod.rs b/rust_snuba/src/pull_consumer/mod.rs new file mode 100644 index 0000000000..2d7fd8f160 --- /dev/null +++ b/rust_snuba/src/pull_consumer/mod.rs @@ -0,0 +1,19 @@ +mod entrypoint; +mod factories; +mod pipelines; + +pub use entrypoint::pull_consumer; + +/// Allowed processors for the fire-and-forget pipeline. +const FIRE_AND_FORGET_PROCESSORS: &[&str] = &[ + "FunctionsMessageProcessor", + "ProfilesMessageProcessor", + "QuerylogProcessor", + "ReplaysProcessor", + "OutcomesProcessor", + "ProfileChunksProcessor", + "LlmProxyCostProcessor", +]; + +/// Allowed processors for the EAP pipeline. +const EAP_PROCESSORS: &[&str] = &["EAPItemsProcessor"]; diff --git a/rust_snuba/src/pull_consumer/pipelines.rs b/rust_snuba/src/pull_consumer/pipelines.rs new file mode 100644 index 0000000000..8d30b63492 --- /dev/null +++ b/rust_snuba/src/pull_consumer/pipelines.rs @@ -0,0 +1,86 @@ +use crate::config; +use crate::config::ProcessorConfig; +use crate::processors::get_cogs_label; +use crate::pull::pipelines::eap::EapPipeline; +use crate::pull::pipelines::fire_and_forget::FireAndForgetPipeline; + +use super::factories::{self, SharedResources}; + +pub(super) fn make_eap_pipeline( + shared: &SharedResources, + consumer_config: &config::ConsumerConfig, + storage: &config::StorageConfig, + processor_config: &ProcessorConfig, + consumer_group: &str, + processing_concurrency: usize, + clickhouse_concurrency: usize, + dry_run: bool, +) -> EapPipeline { + let processor_name = &storage.message_processor.python_class_name; + let source_topic_name = &consumer_config.raw_topic.physical_topic_name; + + assert_eq!(processor_name, "EAPItemsProcessor"); + + let resource_id = + get_cogs_label(processor_name).unwrap_or_else(|| format!("{}_processor", storage.name)); + let (cadence, idle_timeout) = factories::make_flush_timers(consumer_config); + + EapPipeline::new( + factories::make_processor_stage( + crate::processors::eap_items::process_message_row_binary, + processor_config, + ), + processing_concurrency, + factories::make_dlq_handler( + &shared.dlq_producer, + consumer_config.dlq_topic.as_ref(), + dry_run, + ), + factories::make_batch_stage( + consumer_config.max_batch_size as u64, + consumer_config.max_batch_size_calculation, + ), + cadence, + idle_timeout, + factories::make_writer_stage(&shared.ch_writer), + clickhouse_concurrency, + factories::make_commit_log_stage( + &shared.commit_log_producer, + consumer_config.commit_log_topic.as_ref(), + source_topic_name, + consumer_group, + dry_run, + ), + factories::make_cogs_stage( + &shared.cogs_producer, + &consumer_config.accountant_topic, + resource_id, + dry_run, + consumer_config.env.record_cogs, + ), + ) +} + +pub(super) fn make_faf_pipeline( + shared: &SharedResources, + processor: crate::processors::ProcessingFunction, + consumer_config: &config::ConsumerConfig, + processor_config: &ProcessorConfig, + processing_concurrency: usize, + clickhouse_concurrency: usize, +) -> FireAndForgetPipeline { + let (cadence, idle_timeout) = factories::make_flush_timers(consumer_config); + + FireAndForgetPipeline::new( + factories::make_processor_stage(processor, processor_config), + processing_concurrency, + factories::make_batch_stage( + consumer_config.max_batch_size as u64, + consumer_config.max_batch_size_calculation, + ), + cadence, + idle_timeout, + factories::make_writer_stage(&shared.ch_writer), + clickhouse_concurrency, + ) +} From 44abc287dab97c562a991c61d1bdf81fd15375e0 Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Fri, 11 Sep 2026 14:04:08 -0700 Subject: [PATCH 7/7] clippy :smithers: --- rust_snuba/src/pull_consumer/pipelines.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/rust_snuba/src/pull_consumer/pipelines.rs b/rust_snuba/src/pull_consumer/pipelines.rs index 8d30b63492..4839ec38e0 100644 --- a/rust_snuba/src/pull_consumer/pipelines.rs +++ b/rust_snuba/src/pull_consumer/pipelines.rs @@ -6,6 +6,7 @@ use crate::pull::pipelines::fire_and_forget::FireAndForgetPipeline; use super::factories::{self, SharedResources}; +#[allow(clippy::too_many_arguments)] pub(super) fn make_eap_pipeline( shared: &SharedResources, consumer_config: &config::ConsumerConfig,