diff --git a/Cargo.lock b/Cargo.lock index 1ecb861658bbe..d841ec130d968 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2524,6 +2524,7 @@ dependencies = [ "async-trait", "bytes 1.11.1", "chrono", + "criterion", "csv-core", "derivative", "derive_more 2.1.1", diff --git a/lib/codecs/Cargo.toml b/lib/codecs/Cargo.toml index 7cbc0c1c3aaca..08c7bd5abacb3 100644 --- a/lib/codecs/Cargo.toml +++ b/lib/codecs/Cargo.toml @@ -65,6 +65,7 @@ toml.workspace = true similar-asserts = "1.7.0" vector-core = { path = "../vector-core", default-features = false, features = ["vrl", "test"] } rstest = "0.26.1" +criterion = { version = "0.7.0", features = ["html_reports"] } tracing-test = "0.2.6" uuid.workspace = true vrl.workspace = true @@ -74,3 +75,11 @@ arrow = ["dep:arrow"] opentelemetry = ["dep:opentelemetry-proto"] syslog = ["dep:syslog_loose", "dep:strum", "dep:derive_more", "dep:serde-aux", "dep:toml"] test = [] + +[[bench]] +name = "parquet_encode" +harness = false + +[[bench]] +name = "parquet_threads" +harness = false diff --git a/lib/codecs/benches/fixtures/cloudtrail.schema b/lib/codecs/benches/fixtures/cloudtrail.schema new file mode 100644 index 0000000000000..83e1e3bd8e2a3 --- /dev/null +++ b/lib/codecs/benches/fixtures/cloudtrail.schema @@ -0,0 +1,81 @@ +message CLOUDTRAIL { + required int64 exfRecordIndex; + required int64 exfVectorTime (Timestamp(MILLIS, true)); + required binary exfS3Bucket (String); + required binary exfS3ObjectName (String); + + optional binary exfConnectorUid (String); + optional binary exfSourceTags (JSON); + optional binary exfAssumeRoleRequestParameters (String); + optional binary exfAssumeRoleResponseElements (String); + optional binary exfInstancesRequestParameters (String); + optional binary exfInstancesResponseElements (String); + optional binary exfRequestParameters (String); + optional binary exfResponseElements (String); + optional boolean exfRequestParametersOmitted; + optional binary exfS3AdditionalEventData (String); + optional binary exfAdditionalEventData (String); + optional binary exfAdditionalEventDataSessionization (String); + optional binary exfResources (String); + optional binary addendum (String); + optional binary additionalEventData (String); + optional binary annotation (String); + optional binary apiVersion (String); + optional binary awsRegion (String); + optional binary edgeDeviceDetails (String); + optional binary errorCode (String); + optional binary errorMessage (String); + optional binary eventCategory (String); + required binary eventID (String); + optional binary eventName (String); + optional binary eventSource (String); + required int64 eventTime (Timestamp(MILLIS, true)); + optional binary eventType (String); + optional binary eventVersion (String); + optional binary insightDetails (String); + optional boolean managementEvent; + optional boolean readOnly; + optional binary recipientAccountId (String); + optional binary requestID (String); + optional binary requestParameters (String); + optional binary resources (String); + optional binary responseElements (String); + optional binary serviceEventDetails (String); + optional binary sessionCredentialFromConsole (String); + optional binary sharedEventID (String); + optional binary sourceIPAddress (String); + optional binary tlsDetails (String); + optional binary userAgent (String); + optional binary userIdentity.accessKeyId (String); + optional binary userIdentity.accountId (String); + optional binary userIdentity.arn (String); + optional binary userIdentity.credentialId (String); + optional binary userIdentity.identityProvider (String); + optional binary userIdentity.inScopeOf.credentialsIssuedTo (String); + optional binary userIdentity.inScopeOf.issuerType (String); + optional binary userIdentity.inScopeOf.sourceAccount (String); + optional binary userIdentity.inScopeOf.sourceArn (String); + optional binary userIdentity.invokedBy (String); + optional binary userIdentity.invokedByDelegate.accountId (String); + optional binary userIdentity.onBehalfOf.identityStoreArn (String); + optional binary userIdentity.onBehalfOf.userId (String); + optional binary userIdentity.principalId (String); + optional int64 userIdentity.sessionContext.attributes.creationDate (Timestamp(MILLIS, true)); + optional binary userIdentity.sessionContext.attributes.mfaAuthenticated (String); + optional binary userIdentity.sessionContext.ec2IssuedInVPC (String); + optional binary userIdentity.sessionContext.ec2RoleDelivery (String); + optional binary userIdentity.sessionContext.sessionIssuer.accountId (String); + optional binary userIdentity.sessionContext.sessionIssuer.arn (String); + optional binary userIdentity.sessionContext.sessionIssuer.principalId (String); + optional binary userIdentity.sessionContext.sessionIssuer.type (String); + optional binary userIdentity.sessionContext.sessionIssuer.userName (String); + optional boolean userIdentity.sessionContext.assumedRoot; + optional binary userIdentity.sessionContext.signInSessionArn (String); + optional binary userIdentity.sessionContext.sourceIdentity (String); + optional binary userIdentity.sessionContext.webIdFederationData.attributes (String); + optional binary userIdentity.sessionContext.webIdFederationData.federatedProvider (String); + optional binary userIdentity.type (String); + optional binary userIdentity.userName (String); + optional binary vpcEndpointId (String); + optional binary vpcEndpointAccountId (String); +} diff --git a/lib/codecs/benches/parquet_encode.rs b/lib/codecs/benches/parquet_encode.rs new file mode 100644 index 0000000000000..06d29f279685f --- /dev/null +++ b/lib/codecs/benches/parquet_encode.rs @@ -0,0 +1,135 @@ +//! Encodes CloudTrail-shaped events against the real 78-column production schema. +//! +//! The deployed encoder is column-major: `encode()` walks every event once per +//! column, so cost scales with columns x events. This measures that, and gives a +//! baseline to compare any row-major rewrite against. + +use bytes::BytesMut; +use codecs::encoding::ParquetSerializerConfig; +use criterion::{BatchSize, BenchmarkId, Criterion, Throughput, criterion_group, criterion_main}; +use parquet::basic::Compression; +use tokio_util::codec::Encoder; +use vector_core::event::{Event, LogEvent, ObjectMap, Value}; + +/// The production CloudTrail schema: 78 leaf columns, no nested groups. +const SCHEMA: &str = include_str!("fixtures/cloudtrail.schema"); + +/// Every leaf column name in SCHEMA, in declaration order. +/// +/// Names carry literal dots (`userIdentity.accessKeyId`) because cloudtrail.vrl +/// flattens `userIdentity` and merges the flat keys into the record, so the +/// encoder sees top-level keys containing dots -- not a nested object. +fn schema_fields() -> Vec<(String, &'static str)> { + let mut fields = Vec::new(); + for line in SCHEMA.lines() { + let line = line.trim().trim_end_matches(';'); + if line.starts_with("message") || line.starts_with('}') || line.is_empty() { + continue; + } + let mut parts = line.split_whitespace(); + let _optionality = parts.next(); + let phys = match parts.next() { + Some(p) => p, + None => continue, + }; + let name = match parts.next() { + Some(n) => n.to_string(), + None => continue, + }; + // Timestamp-annotated int64 columns still take an integer value. + let kind = match phys { + "binary" => "binary", + "int64" => "int64", + "boolean" => "boolean", + other => panic!("unhandled physical type in schema: {other}"), + }; + fields.push((name, kind)); + } + fields +} + +/// Distinct values a column takes across the corpus. +/// +/// Cardinality matters as much as length here: parquet dictionary-encodes each +/// column, and a corpus of unique-per-row strings inflates dictionary growth +/// with row count, which would masquerade as encoder scaling. These approximate +/// real CloudTrail -- a handful of regions, a few hundred sources, IDs unique. +fn cardinality(name: &str) -> usize { + match name { + "awsRegion" => 30, + "eventSource" => 200, + "eventName" => 2_000, + "eventType" | "eventCategory" | "eventVersion" | "userIdentity.type" => 8, + "recipientAccountId" | "userIdentity.accountId" => 40, + "userAgent" => 150, + "sourceIPAddress" => 5_000, + // Identifiers are genuinely unique per record. + "eventID" | "requestID" | "sharedEventID" | "exfS3ObjectName" => usize::MAX, + _ => 500, + } +} + +/// Builds one event with every column populated. String lengths are in the range +/// real CloudTrail records occupy -- ARNs and user agents dominate the payload. +fn make_event(fields: &[(String, &'static str)], i: usize) -> Event { + let mut map = ObjectMap::new(); + for (name, kind) in fields { + let value = match *kind { + "int64" => Value::from(1_700_000_000_000i64 + i as i64), + "boolean" => Value::from(i % 2 == 0), + _ => { + let card = cardinality(name); + let n = if card == usize::MAX { i } else { i % card }; + match name.as_str() { + x if x.ends_with("arn") || x.ends_with("Arn") => Value::from(format!( + "arn:aws:sts::123456789012:assumed-role/ExampleRoleName/session-{n}" + )), + "userAgent" => Value::from(format!( + "aws-sdk-go/1.44.{n} (go1.19.3; linux; amd64) exec-env/AWS_Lambda_go1.x" + )), + "sourceIPAddress" => { + Value::from(format!("10.{}.{}.{}", n / 65536 % 256, n / 256 % 256, n % 256)) + } + _ => Value::from(format!("{name}-value-{n}")), + } + } + }; + map.insert(name.as_str().into(), value); + } + Event::Log(LogEvent::from(map)) +} + +fn bench_encode(c: &mut Criterion) { + let fields = schema_fields(); + assert_eq!(fields.len(), 78, "schema fixture drifted from 78 columns"); + + let mut group = c.benchmark_group("parquet_encode_cloudtrail"); + group.sample_size(10); + + for batch in [1_000usize, 10_000, 50_000] { + let events: Vec = (0..batch).map(|i| make_event(&fields, i)).collect(); + group.throughput(Throughput::Elements(batch as u64)); + group.bench_with_input(BenchmarkId::from_parameter(batch), &events, |b, events| { + b.iter_batched( + || { + ( + ParquetSerializerConfig::new(SCHEMA.to_string()) + .build(Compression::UNCOMPRESSED) + .expect("schema builds"), + events.clone(), + BytesMut::with_capacity(16 * 1024 * 1024), + ) + }, + |(mut ser, events, mut buf)| { + ser.encode(events, &mut buf).expect("encode succeeds"); + buf.len() + }, + BatchSize::LargeInput, + ); + }); + } + group.finish(); +} + +criterion_group!(benches, bench_encode); +criterion_main!(benches); diff --git a/lib/codecs/benches/parquet_threads.rs b/lib/codecs/benches/parquet_threads.rs new file mode 100644 index 0000000000000..baae26557ca37 --- /dev/null +++ b/lib/codecs/benches/parquet_threads.rs @@ -0,0 +1,150 @@ +//! Thread-scaling measurement for the parquet codec. +//! +//! The production failure was not raw encode cost -- it was 32 workers contending +//! on the refcount of ONE shared allocation. Events parsed out of a single S3 +//! object hold `Bytes` that all slice the same buffer, so every clone/drop during +//! encoding hits the same cache line. +//! +//! Criterion's per-iteration `format!` corpus cannot show that: each value gets +//! its own allocation and therefore its own refcount. Here every string value is +//! a slice of one shared blob, which is what production looks like. +//! +//! Run with the codec's ROW_GROUP_ROWS set to usize::MAX and again at 1024 to +//! compare. Prints aggregate events/s per thread count. + +use bytes::{Bytes, BytesMut}; +use codecs::encoding::ParquetSerializerConfig; +use parquet::basic::Compression; +use std::time::Instant; +use tokio_util::codec::Encoder; +use vector_core::event::{Event, LogEvent, ObjectMap, Value}; + +const SCHEMA: &str = include_str!("fixtures/cloudtrail.schema"); +const BATCH: usize = 10_000; + +fn schema_fields() -> Vec<(String, &'static str)> { + let mut fields = Vec::new(); + for line in SCHEMA.lines() { + let line = line.trim().trim_end_matches(';'); + if line.starts_with("message") || line.starts_with('}') || line.is_empty() { + continue; + } + let mut parts = line.split_whitespace(); + let _optionality = parts.next(); + let phys = parts.next().unwrap_or(""); + let name = match parts.next() { + Some(n) => n.to_string(), + None => continue, + }; + let kind = match phys { + "binary" => "binary", + "int64" => "int64", + "boolean" => "boolean", + _ => continue, + }; + fields.push((name, kind)); + } + fields +} + +/// One allocation holding every distinct string the corpus uses, mirroring a +/// parsed S3 object. Returns the blob plus the (offset, len) of each entry. +fn shared_blob(fields: &[(String, &'static str)], rows: usize) -> (Bytes, Vec>) { + let mut buf = BytesMut::with_capacity(8 * 1024 * 1024); + let mut spans: Vec> = Vec::with_capacity(fields.len()); + + for (name, kind) in fields { + let mut per_field = Vec::new(); + if *kind != "binary" { + spans.push(per_field); + continue; + } + // Distinct values per column, as in the criterion corpus. + let card = match name.as_str() { + "awsRegion" => 30, + "eventSource" => 200, + "eventName" => 2_000, + "eventID" | "requestID" | "sharedEventID" | "exfS3ObjectName" => rows, + _ => 500, + } + .min(rows); + for n in 0..card { + let s = if name.ends_with("arn") || name.ends_with("Arn") { + format!("arn:aws:sts::123456789012:assumed-role/ExampleRoleName/session-{n}") + } else { + format!("{name}-value-{n}") + }; + let off = buf.len(); + buf.extend_from_slice(s.as_bytes()); + per_field.push((off, s.len())); + } + spans.push(per_field); + } + (buf.freeze(), spans) +} + +fn make_events( + fields: &[(String, &'static str)], + blob: &Bytes, + spans: &[Vec<(usize, usize)>], + rows: usize, +) -> Vec { + (0..rows) + .map(|i| { + let mut map = ObjectMap::new(); + for (fi, (name, kind)) in fields.iter().enumerate() { + let value = match *kind { + "int64" => Value::from(1_700_000_000_000i64 + i as i64), + "boolean" => Value::from(i % 2 == 0), + _ => { + let per_field = &spans[fi]; + let (off, len) = per_field[i % per_field.len()]; + // Slice of the shared blob: same refcount for every event. + Value::Bytes(blob.slice(off..off + len)) + } + }; + map.insert(name.as_str().into(), value); + } + Event::Log(LogEvent::from(map)) + }) + .collect() +} + +fn main() { + let fields = schema_fields(); + let (blob, spans) = shared_blob(&fields, BATCH); + println!( + "schema columns: {} batch: {} shared blob: {:.1} MiB", + fields.len(), + BATCH, + blob.len() as f64 / 1048576.0 + ); + println!("{:>7} {:>14} {:>12}", "threads", "events/s", "vs 1 thread"); + + let mut single = 0.0f64; + for threads in [1usize, 2, 4, 8, 12] { + // Each thread gets its own event vec, but all values share ONE blob. + let per_thread: Vec> = (0..threads) + .map(|_| make_events(&fields, &blob, &spans, BATCH)) + .collect(); + + let start = Instant::now(); + std::thread::scope(|scope| { + for events in &per_thread { + scope.spawn(|| { + let mut ser = ParquetSerializerConfig::new(SCHEMA.to_string()) + .build(Compression::UNCOMPRESSED) + .expect("schema builds"); + let mut buf = BytesMut::with_capacity(16 * 1024 * 1024); + ser.encode(events.clone(), &mut buf).expect("encode"); + }); + } + }); + let elapsed = start.elapsed().as_secs_f64(); + let total = (threads * BATCH) as f64 / elapsed; + if threads == 1 { + single = total; + } + println!("{threads:>7} {total:>14.0} {:>11.2}x", total / single); + } +} diff --git a/lib/codecs/src/encoding/format/parquet.rs b/lib/codecs/src/encoding/format/parquet.rs index 3f00e476967b9..a521668d15c14 100644 --- a/lib/codecs/src/encoding/format/parquet.rs +++ b/lib/codecs/src/encoding/format/parquet.rs @@ -294,6 +294,9 @@ impl ParquetSerializer { } } +/// Rows per parquet row group. +const ROW_GROUP_ROWS: usize = 8192; + impl Encoder> for ParquetSerializer { type Error = vector_common::Error; @@ -315,72 +318,82 @@ impl Encoder> for ParquetSerializer { let mut parquet_writer = SerializedFileWriter::new(buffer.writer(), self.schema.clone(), props)?; - let mut row_group_writer = parquet_writer.next_row_group()?; - while let Some(mut column_writer) = row_group_writer.next_column()? { - match column_writer.untyped() { - BoolColumnWriter(writer) => { - let desc = writer.get_descriptor().clone(); - self.process( - &events, - &desc, - |value| match value { - Value::Boolean(value) => Ok(*value), - _ => Err(ParquetSerializerError::invalid_type( - &desc, value, "boolean", - )), - }, - writer, - )? - } - Int64ColumnWriter(writer) => { - let desc = writer.get_descriptor().clone(); - self.process( - &events, - &desc, - |value| match value { - Value::Integer(value) => Ok(*value), - _ => Err(ParquetSerializerError::invalid_type( - &desc, value, "integer", - )), - }, - writer, - )? - } - DoubleColumnWriter(writer) => { - let desc = writer.get_descriptor().clone(); - self.process( - &events, - &desc, - |value| match value { - Value::Float(value) => Ok(value.into_inner()), - _ => Err(ParquetSerializerError::invalid_type(&desc, value, "float")), - }, - writer, - )? - } - ByteArrayColumnWriter(writer) => { - let desc = writer.get_descriptor().clone(); - self.process( - &events, - &desc, - |value| match value { - Value::Bytes(value) => Ok(value.clone().into()), - _ => Err(ParquetSerializerError::invalid_type(&desc, value, "string")), - }, - writer, - )? - } - FixedLenByteArrayColumnWriter(_) => { - panic!("Fixed len byte array is not supported."); + // One row group per chunk. Each column pass re-walks its slice of events, + // so an unbounded batch makes all 78 passes stream from memory instead of + // cache; chunking keeps the working set resident. Memory is unchanged -- + // still one column buffered at a time. + for events in events.chunks(ROW_GROUP_ROWS) { + let mut row_group_writer = parquet_writer.next_row_group()?; + while let Some(mut column_writer) = row_group_writer.next_column()? { + match column_writer.untyped() { + BoolColumnWriter(writer) => { + let desc = writer.get_descriptor().clone(); + self.process( + &events, + &desc, + |value| match value { + Value::Boolean(value) => Ok(*value), + _ => Err(ParquetSerializerError::invalid_type( + &desc, value, "boolean", + )), + }, + writer, + )? + } + Int64ColumnWriter(writer) => { + let desc = writer.get_descriptor().clone(); + self.process( + &events, + &desc, + |value| match value { + Value::Integer(value) => Ok(*value), + _ => Err(ParquetSerializerError::invalid_type( + &desc, value, "integer", + )), + }, + writer, + )? + } + DoubleColumnWriter(writer) => { + let desc = writer.get_descriptor().clone(); + self.process( + &events, + &desc, + |value| match value { + Value::Float(value) => Ok(value.into_inner()), + _ => { + Err(ParquetSerializerError::invalid_type(&desc, value, "float")) + } + }, + writer, + )? + } + ByteArrayColumnWriter(writer) => { + let desc = writer.get_descriptor().clone(); + self.process( + &events, + &desc, + |value| match value { + Value::Bytes(value) => Ok(value.clone().into()), + _ => Err(ParquetSerializerError::invalid_type( + &desc, value, "string", + )), + }, + writer, + )? + } + FixedLenByteArrayColumnWriter(_) => { + panic!("Fixed len byte array is not supported."); + } + Int32ColumnWriter(_) => panic!("Int32 is not supported."), + Int96ColumnWriter(_) => panic!("Int96 is not supported."), + FloatColumnWriter(_) => panic!("Float32 is not supported."), } - Int32ColumnWriter(_) => panic!("Int32 is not supported."), - Int96ColumnWriter(_) => panic!("Int96 is not supported."), - FloatColumnWriter(_) => panic!("Float32 is not supported."), + column_writer.close()?; } - column_writer.close()?; - } - row_group_writer.close()?; + row_group_writer.close()?; + } parquet_writer.close()?; Ok(()) @@ -447,12 +460,8 @@ impl<'a, T, F: Fn(&Value) -> Result> Column<'a, T, F> fn extract_column(&mut self, events: &[Event]) -> Result<(), ParquetSerializerError> { for event in events { let res = match event { - Event::Log(log) => { - self.extract_value(log.value(), Level::root()) - } - Event::Trace(trace) => { - self.extract_value(trace.value(), Level::root()) - } + Event::Log(log) => self.extract_value(log.value(), Level::root()), + Event::Trace(trace) => self.extract_value(trace.value(), Level::root()), Event::Metric(_) => { panic!("Metrics are not supported."); }