#![allow(clippy::print_stdout, clippy::print_stderr)]
use serde::{Deserialize, Serialize};
use spate::avro::AvroDeserializerBuilder;
use spate::clickhouse::{NativeEncoder, ShardKey};
use spate::kafka::KafkaSource;
use spate::prelude::*;
use std::path::Path;
#[derive(Debug, Deserialize)]
enum StorefrontEvent {
OrderPlaced(serde::de::IgnoredAny),
PaymentCaptured(PaymentRow),
RefundIssued(RefundRow),
}
#[derive(Debug, Deserialize, Serialize)]
struct PaymentRow {
order_id: u64,
amount_cents: u64,
}
#[derive(Debug, Deserialize, Serialize)]
struct RefundRow {
order_id: u64,
amount_cents: u64,
reason: String,
}
fn payment_key(row: &PaymentRow) -> ShardKey<'_> {
ShardKey::U64(row.order_id)
}
fn refund_key(row: &RefundRow) -> ShardKey<'_> {
ShardKey::U64(row.order_id)
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
let config_path = std::env::var("SPATE_CONFIG")
.unwrap_or_else(|_| "crates/spate/examples/multi_table_split.yaml".to_string());
let pipeline = Pipeline::from_path(Path::new(&config_path))?;
let source = KafkaSource::from_component_config(&pipeline.config().source)?;
let deser_section = pipeline
.config()
.deserializer
.as_ref()
.ok_or("this pipeline requires a `deserializer` section")?;
let deserializer =
AvroDeserializerBuilder::from_component(deser_section, &pipeline.io_handle())?
.build_serde::<StorefrontEvent>()?;
let payments_sink = spate::clickhouse::config::from_component_config(
pipeline.config().sink_config("payments")?,
)?;
let refunds_sink = spate::clickhouse::config::from_component_config(
pipeline.config().sink_config("refunds")?,
)?;
let payments_router = payments_sink.router::<Owned<PaymentRow>>(payment_key);
let refunds_router = refunds_sink.router::<Owned<RefundRow>>(refund_key);
let payments_enc =
NativeEncoder::<Owned<PaymentRow>>::new(pipeline.block_on(payments_sink.native_schema())?);
let refunds_enc =
NativeEncoder::<Owned<RefundRow>>::new(pipeline.block_on(refunds_sink.native_schema())?);
let report = pipeline
.add_sink("payments", payments_sink)?
.add_sink("refunds", refunds_sink)?
.chains(move |ctx| {
let mut split = chain::<Owned<StorefrontEvent>, _>(deserializer.clone())
.with_metrics(ctx.pipeline.clone(), "main")
.split(ErrorPolicy::Skip);
let payments = split.add::<Owned<PaymentRow>, _, _>(
payments_enc.clone(),
payments_router.clone(),
ctx.sink("payments"),
);
let refunds = split.add::<Owned<RefundRow>, _, _>(
refunds_enc.clone(),
refunds_router.clone(),
ctx.sink("refunds"),
);
split
.route(move |event: StorefrontEvent, out| match event {
StorefrontEvent::PaymentCaptured(row) => out.emit(payments, row),
StorefrontEvent::RefundIssued(row) => out.emit(refunds, row),
StorefrontEvent::OrderPlaced(_) => {}
})
.build()
})
.run(source)?;
report.log();
std::process::exit(report.exit_code());
}