Skip to content
Open
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
44 changes: 42 additions & 2 deletions quickwit/quickwit-ingest/benches/mrecord_bench.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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];
Expand Down Expand Up @@ -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<Bytes> = (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);
9 changes: 9 additions & 0 deletions quickwit/quickwit-ingest/src/ingest_v2/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
79 changes: 67 additions & 12 deletions quickwit/quickwit-serve/src/ingest_api/rest_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -193,14 +193,7 @@ async fn ingest_v2(
ingest_options: IngestOptions,
ingest_router: IngestRouterServiceClient,
) -> Result<RestIngestResponse, IngestServiceError> {
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();
Expand Down Expand Up @@ -279,6 +272,51 @@ pub(crate) fn lines(body: &Bytes) -> impl Iterator<Item = &[u8]> {
.filter(|line| !is_empty_or_blank_line(line))
}

fn build_doc_batch_v2_from_ndjson_body(doc_buffer: Bytes) -> Option<DocBatchV2> {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
fn build_doc_batch_v2_from_ndjson_body(doc_buffer: Bytes) -> Option<DocBatchV2> {
fn build_doc_batch_v2_from_ndjson_body(ndjson_body: Bytes) -> Option<DocBatchV2> {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Any chance we can abstract away the splitting and parsing part so we can reused in Pomsky for other handlers?

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())
Expand All @@ -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]
Expand All @@ -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<Bytes> = 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,
Expand Down
Loading