Skip to content
Closed
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
1 change: 1 addition & 0 deletions Cargo.lock

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

9 changes: 9 additions & 0 deletions lib/codecs/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
81 changes: 81 additions & 0 deletions lib/codecs/benches/fixtures/cloudtrail.schema
Original file line number Diff line number Diff line change
@@ -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);
}
135 changes: 135 additions & 0 deletions lib/codecs/benches/parquet_encode.rs
Original file line number Diff line number Diff line change
@@ -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<Event> = (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);
150 changes: 150 additions & 0 deletions lib/codecs/benches/parquet_threads.rs
Original file line number Diff line number Diff line change
@@ -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<Vec<(usize, usize)>>) {
let mut buf = BytesMut::with_capacity(8 * 1024 * 1024);
let mut spans: Vec<Vec<(usize, usize)>> = 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<Event> {
(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<Vec<Event>> = (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);
}
}
Loading
Loading