diff --git a/quickwit/quickwit-ingest/benches/mrecord_bench.rs b/quickwit/quickwit-ingest/benches/mrecord_bench.rs index 228a5567877..5e750d3d914 100644 --- a/quickwit/quickwit-ingest/benches/mrecord_bench.rs +++ b/quickwit/quickwit-ingest/benches/mrecord_bench.rs @@ -28,7 +28,8 @@ use std::time::{Duration, Instant}; use bytes::{BufMut, Bytes}; use criterion::{BenchmarkId, Criterion, Throughput, criterion_group, criterion_main}; use mrecordlog::MultiRecordLog; -use quickwit_ingest::MRecord; +use quickwit_ingest::{DocBatchV2Builder, MRecord}; +use quickwit_proto::types::DocUidGenerator; /// Representative document sizes (bytes): a small log line up to a large structured event. const DOC_SIZES: [usize; 4] = [128, 1_024, 8_192, 65_536]; @@ -180,5 +181,44 @@ fn bench_append_batch(criterion: &mut Criterion) { group.finish(); } -criterion_group!(benches, bench_encode, bench_decode, bench_append_batch); +fn bench_doc_batch_builder(criterion: &mut Criterion) { + let mut group = criterion.benchmark_group("doc_batch_v2_builder"); + let docs: Vec = (0..10_000) + .map(|index| make_doc(BATCH_DOC_SIZES[index % BATCH_DOC_SIZES.len()])) + .collect(); + let total_bytes: usize = docs.iter().map(Bytes::len).sum(); + group.throughput(Throughput::Bytes(total_bytes as u64)); + + group.bench_function("default", |bencher| { + bencher.iter(|| { + let mut builder = DocBatchV2Builder::default(); + let mut doc_uid_generator = DocUidGenerator::default(); + for doc in &docs { + builder.add_doc(doc_uid_generator.next_doc_uid(), doc); + } + black_box(builder.build()) + }); + }); + + group.bench_function("preallocated", |bencher| { + bencher.iter(|| { + let mut builder = DocBatchV2Builder::with_capacity(total_bytes); + let mut doc_uid_generator = DocUidGenerator::default(); + for doc in &docs { + builder.add_doc(doc_uid_generator.next_doc_uid(), doc); + } + black_box(builder.build()) + }); + }); + + group.finish(); +} + +criterion_group!( + benches, + bench_encode, + bench_decode, + bench_append_batch, + bench_doc_batch_builder +); criterion_main!(benches); diff --git a/quickwit/quickwit-ingest/src/ingest_v2/mod.rs b/quickwit/quickwit-ingest/src/ingest_v2/mod.rs index 5a6c6e2c791..e77a3dc358d 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/mod.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/mod.rs @@ -153,6 +153,15 @@ pub struct DocBatchV2Builder { } impl DocBatchV2Builder { + /// Creates a builder with enough room for `doc_buffer_capacity` bytes of document payloads. + pub fn with_capacity(doc_buffer_capacity: usize) -> Self { + Self { + doc_uids: Vec::new(), + doc_buffer: BytesMut::with_capacity(doc_buffer_capacity), + doc_lengths: Vec::new(), + } + } + /// Adds a document to the batch. pub fn add_doc(&mut self, doc_uid: DocUid, doc: &[u8]) { self.doc_uids.push(doc_uid); diff --git a/quickwit/quickwit-serve/src/ingest_api/rest_handler.rs b/quickwit/quickwit-serve/src/ingest_api/rest_handler.rs index 9040c8fd700..d3e48d447cd 100644 --- a/quickwit/quickwit-serve/src/ingest_api/rest_handler.rs +++ b/quickwit/quickwit-serve/src/ingest_api/rest_handler.rs @@ -15,13 +15,13 @@ use bytes::{Buf, Bytes}; use quickwit_config::{INGEST_V2_SOURCE_ID, IngestApiConfig, validate_identifier}; use quickwit_ingest::{ - CommitType, DocBatchBuilder, DocBatchV2Builder, FetchResponse, IngestRequest, IngestService, - IngestServiceClient, IngestServiceError, TailRequest, + CommitType, DocBatchBuilder, FetchResponse, IngestRequest, IngestService, IngestServiceClient, + IngestServiceError, TailRequest, }; -use quickwit_proto::ingest::CommitTypeV2; use quickwit_proto::ingest::router::{ IngestRequestV2, IngestRouterService, IngestRouterServiceClient, IngestSubrequest, }; +use quickwit_proto::ingest::{CommitTypeV2, DocBatchV2}; use quickwit_proto::types::{DocUidGenerator, IndexId}; use serde::Deserialize; use warp::{Filter, Rejection}; @@ -193,14 +193,7 @@ async fn ingest_v2( ingest_options: IngestOptions, ingest_router: IngestRouterServiceClient, ) -> Result { - let mut doc_batch_builder = DocBatchV2Builder::default(); - let mut doc_uid_generator = DocUidGenerator::default(); - - for doc in lines(&body.content) { - doc_batch_builder.add_doc(doc_uid_generator.next_doc_uid(), doc); - } - drop(body); - let doc_batch_opt = doc_batch_builder.build(); + let doc_batch_opt = build_doc_batch_v2_from_ndjson_body(body.content); let Some(doc_batch) = doc_batch_opt else { let response = RestIngestResponse::default(); @@ -279,6 +272,51 @@ pub(crate) fn lines(body: &Bytes) -> impl Iterator { .filter(|line| !is_empty_or_blank_line(line)) } +fn build_doc_batch_v2_from_ndjson_body(doc_buffer: Bytes) -> Option { + let mut doc_uids = Vec::new(); + let mut doc_lengths = Vec::new(); + let mut doc_uid_generator = DocUidGenerator::default(); + let mut segment_start = 0usize; + let mut line_start = 0usize; + + for (position, byte) in doc_buffer.iter().enumerate() { + if *byte != b'\n' { + continue; + } + let line = &doc_buffer[line_start..position]; + if !is_empty_or_blank_line(line) { + doc_uids.push(doc_uid_generator.next_doc_uid()); + doc_lengths.push((position + 1 - segment_start) as u32); + segment_start = position + 1; + } + line_start = position + 1; + } + + let line = &doc_buffer[line_start..]; + if !is_empty_or_blank_line(line) { + doc_uids.push(doc_uid_generator.next_doc_uid()); + doc_lengths.push((doc_buffer.len() - segment_start) as u32); + segment_start = doc_buffer.len(); + } + + if doc_uids.is_empty() { + return None; + } + if segment_start < doc_buffer.len() { + let trailing_whitespace_len = doc_buffer.len() - segment_start; + let last_doc_len = doc_lengths + .last_mut() + .expect("doc lengths should not be empty"); + *last_doc_len += trailing_whitespace_len as u32; + } + + Some(DocBatchV2 { + doc_uids, + doc_buffer, + doc_lengths, + }) +} + #[inline] fn is_empty_or_blank_line(line: &[u8]) -> bool { line.is_empty() || line.iter().all(|ch| ch.is_ascii_whitespace()) @@ -298,7 +336,7 @@ pub(crate) mod tests { }; use quickwit_proto::ingest::router::IngestRouterServiceClient; - use super::{RestIngestResponse, ingest_api_handlers}; + use super::{RestIngestResponse, build_doc_batch_v2_from_ndjson_body, ingest_api_handlers}; use crate::ingest_api::lines; #[test] @@ -319,6 +357,23 @@ pub(crate) mod tests { } } + #[test] + fn test_build_doc_batch_v2_from_ndjson_body_zero_copy() { + let body = Bytes::from_static(b"\n {\"id\":1}\n\n{\"id\":2}\n \n"); + let doc_batch = build_doc_batch_v2_from_ndjson_body(body.clone()).unwrap(); + assert_eq!(doc_batch.num_docs(), 2); + assert_eq!(doc_batch.doc_buffer, body); + + let docs: Vec = doc_batch.docs().map(|(_doc_uid, doc)| doc).collect(); + assert_eq!(str::from_utf8(&docs[0]).unwrap(), "\n {\"id\":1}\n"); + assert_eq!(str::from_utf8(&docs[1]).unwrap(), "\n{\"id\":2}\n \n"); + } + + #[test] + fn test_build_doc_batch_v2_from_blank_body() { + assert!(build_doc_batch_v2_from_ndjson_body(Bytes::from_static(b"\n \n\t")).is_none()); + } + pub(crate) async fn setup_ingest_v1_service( queues: &[&str], config: &IngestApiConfig,