Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions rust_snuba/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 4 additions & 2 deletions rust_snuba/src/accepted_outcomes_consumer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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 {
Expand Down
6 changes: 4 additions & 2 deletions rust_snuba/src/consumer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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),
Expand Down
3 changes: 2 additions & 1 deletion rust_snuba/src/factory_v2.rs
Original file line number Diff line number Diff line change
Expand Up @@ -246,7 +246,8 @@ impl ProcessingStrategyFactory<KafkaPayload> 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,
Expand Down
20 changes: 10 additions & 10 deletions rust_snuba/src/pull/pipelines/eap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -64,16 +64,16 @@ impl EapPipeline {
impl Pipeline for EapPipeline {
type Output = BatchMetadata;

fn stream<'a>(
&'a self,
source: impl Stream<Item = StageResult<KafkaPayload>> + 'a,
) -> impl Stream<Item = StageResult<BatchMetadata>> + 'a {
fn stream(
self,
source: impl Stream<Item = StageResult<KafkaPayload>> + Send,
) -> impl Stream<Item = StageResult<BatchMetadata>> + 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)
Comment thread
tryangul marked this conversation as resolved.
.apply_concurrent(self.writer, self.writer_concurrency)
.apply(self.commit_log)
.apply(self.cogs)
.on_reject(self.dlq_handler)
}
}
16 changes: 8 additions & 8 deletions rust_snuba/src/pull/pipelines/fire_and_forget.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,14 +56,14 @@ impl FireAndForgetPipeline {
impl Pipeline for FireAndForgetPipeline {
type Output = BatchMetadata;

fn stream<'a>(
&'a self,
source: impl Stream<Item = StageResult<KafkaPayload>> + 'a,
) -> impl Stream<Item = StageResult<BatchMetadata>> + 'a {
fn stream(
self,
source: impl Stream<Item = StageResult<KafkaPayload>> + Send,
) -> impl Stream<Item = StageResult<BatchMetadata>> + 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)
}
}
9 changes: 4 additions & 5 deletions rust_snuba/src/pull/stages/clickhouse_writer_stage.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
use std::sync::Arc;
use std::time::Instant;

use chrono::Utc;
Expand All @@ -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<dyn ClickHouseWriter>,
writer: Arc<dyn ClickHouseWriter>,
}

impl ClickHouseWriterStage {
pub fn new(writer: impl ClickHouseWriter + 'static) -> Self {
Self {
writer: Box::new(writer),
}
pub fn new(writer: Arc<dyn ClickHouseWriter>) -> Self {
Self { writer }
}
}

Expand Down
4 changes: 2 additions & 2 deletions rust_snuba/src/pull/stages/cogs_stage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,12 +37,12 @@ pub struct CogsStage {

impl CogsStage {
pub fn new(
producer: impl Producer<KafkaPayload> + 'static,
producer: Arc<dyn Producer<KafkaPayload>>,
destination: Topic,
resource_id: String,
) -> Self {
Self {
producer: Arc::new(producer),
producer,
destination: TopicOrPartition::Topic(destination),
resource_id,
logged_warning: AtomicBool::new(false),
Expand Down
4 changes: 2 additions & 2 deletions rust_snuba/src/pull/stages/commit_log_stage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,13 +31,13 @@ pub struct CommitLogStage {

impl CommitLogStage {
pub fn new(
producer: impl Producer<KafkaPayload> + 'static,
producer: Arc<dyn Producer<KafkaPayload>>,
destination: Topic,
source_topic: Topic,
consumer_group: String,
) -> Self {
Self {
producer: Arc::new(producer),
producer,
destination: TopicOrPartition::Topic(destination),
source_topic,
consumer_group,
Expand Down
3 changes: 1 addition & 2 deletions rust_snuba/src/pull/test_fixtures/sources.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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();

Expand Down
8 changes: 4 additions & 4 deletions rust_snuba/src/pull/tests/eap_pipeline_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
),
Expand Down
2 changes: 1 addition & 1 deletion rust_snuba/src/pull/tests/processor_pipeline_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
);

Expand Down
6 changes: 0 additions & 6 deletions rust_snuba/src/pull/writer/clickhouse_writer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,3 @@ use futures::future::BoxFuture;
pub trait ClickHouseWriter: Send + Sync {
fn write(&self, body: Vec<u8>) -> BoxFuture<'_, anyhow::Result<()>>;
}

impl<T: ClickHouseWriter> ClickHouseWriter for std::sync::Arc<T> {
fn write(&self, body: Vec<u8>) -> BoxFuture<'_, anyhow::Result<()>> {
(**self).write(body)
}
}
Loading
Loading