Repository navigation
Track ingester shard queue size - #6854
nadav-govari wants to merge 2 commits into
Conversation
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: cca58bdad8
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| pub(super) fn doc_batch_size(doc_batch: &DocBatchV2, force_commit: bool) -> ByteSize { | ||
| let estimated_size = estimate_size(doc_batch); |
There was a problem hiding this comment.
Count only bytes actually appended to the queue
estimate_size includes DocBatchV2::num_bytes(), which already adds four bytes for every doc_lengths entry, but those length-table bytes are not written to the WAL. Consequently, the two 12-byte documents in the new test produce 36 bytes here while read_queue_size observes 28, so the test fails; in production, active shards also accumulate an extra four bytes per document, and truncation subtracts only actual evicted bytes, leaving a permanently inflated size that will distort the planned routing logic. Compute this from the document-buffer length plus each encoded mrecord header, adding the commit header only when required.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Codex is right... What we want is:
if doc_batch.is_empty() {
return ByteSize::b(0);
}
let num_records = doc_batch.num_docs() + usize::from(force_commit);
let num_bytes = doc_batch.doc_buffer.len() + num_records * MRECORD_HEADER_LEN;
ByteSize::b(num_bytes as u64)This PR should add a unit test that adds a bunch of records, truncates some, and assert that the queue_size doesn't grow overtime because the number of bytes estimated by us vs. mrecordlog disagree.
| .checked_sub(1) | ||
| .map(Position::offset) | ||
| .unwrap_or(Position::Beginning); | ||
| let queue_size = read_queue_size(&mrecordlog, &queue_id); |
There was a problem hiding this comment.
Keep WAL-size recovery off the Tokio worker
On restart with a large retained WAL, this call synchronously walks every record in every queue to sum payload lengths. init runs in the future spawned by IngesterState::load, so the scan monopolizes a Tokio worker and delays the ingester reaching Ready, with cost proportional to the total retained record count; recover the sizes from maintained WAL metadata or move this CPU-intensive scan off the async executor.
AGENTS.md reference: AGENTS.md:L87-L87
Useful? React with 👍 / 👎.
guilload
left a comment
There was a problem hiding this comment.
This makes me want to delete DocBatchV2::num_bytes() too. Can you look into it?
| let queue_id = subrequest.queue_id; | ||
|
|
||
| let batch_num_docs = subrequest.doc_batch.num_docs() as u64; | ||
| let batch_size = doc_batch_size(&subrequest.doc_batch, force_commit); |
There was a problem hiding this comment.
We should consolidate the two functions doc_batch_size and estimate_size into a single one instead of having one calling the other.
There was a problem hiding this comment.
Makes sense, and since we're not really estimating, as we know the size, it's just doc_batch_size now with the implementation you suggested.
| pub(super) fn doc_batch_size(doc_batch: &DocBatchV2, force_commit: bool) -> ByteSize { | ||
| let estimated_size = estimate_size(doc_batch); |
There was a problem hiding this comment.
Codex is right... What we want is:
if doc_batch.is_empty() {
return ByteSize::b(0);
}
let num_records = doc_batch.num_docs() + usize::from(force_commit);
let num_bytes = doc_batch.doc_buffer.len() + num_records * MRECORD_HEADER_LEN;
ByteSize::b(num_bytes as u64)This PR should add a unit test that adds a bunch of records, truncates some, and assert that the queue_size doesn't grow overtime because the number of bytes estimated by us vs. mrecordlog disagree.
| WAL_BYTES_WRITTEN_TRUNCATE.inc_by(outcome.wal_bytes_written); | ||
| }) | ||
| .map(|outcome| outcome.evicted_records) | ||
| .map(|outcome| ByteSize::b(outcome.evicted_bytes as u64)) |
cca58bd to
a906e6e
Compare
Description
Record the size of the queue on the shard. Will be used for logic to pick which shard to route requests to.
How was this PR tested?
Unit tests.