Skip to content

Track ingester shard queue size - #6854

Open
nadav-govari wants to merge 2 commits into
nadav/no-gossipfrom
feature/rework-scaler
Open

nadav-govari wants to merge 2 commits into
nadav/no-gossipfrom
feature/rework-scaler

Conversation

@nadav-govari

@nadav-govari nadav-govari commented Sep 30, 2026 •

Copy link
Copy Markdown
Collaborator

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.

@nadav-govari
nadav-govari added this pull request to stack #6855 September 30, 2026 15:22
@nadav-govari
nadav-govari requested a review from a team as a code owner September 30, 2026 15:22
@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Sep 30, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-10-07T14:46:57.240730Z 1e17f3b New commits
ℹ️ 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" or "@codex security review".

Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 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".

Comment on lines +102 to +103
pub(super) fn doc_batch_size(doc_batch: &DocBatchV2, force_commit: bool) -> ByteSize {
let estimated_size = estimate_size(doc_batch);

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge 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 👍 / 👎.

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.

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);

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge 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 guilload left a comment

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.

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);

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.

We should consolidate the two functions doc_batch_size and estimate_size into a single one instead of having one calling the other.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

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.

Comment on lines +102 to +103
pub(super) fn doc_batch_size(doc_batch: &DocBatchV2, force_commit: bool) -> ByteSize {
let estimated_size = estimate_size(doc_batch);

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.

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.

Comment thread quickwit/quickwit-ingest/src/ingest_v2/mrecordlog_utils.rs Outdated
Comment thread quickwit/quickwit-ingest/src/ingest_v2/state.rs Outdated
WAL_BYTES_WRITTEN_TRUNCATE.inc_by(outcome.wal_bytes_written);
})
.map(|outcome| outcome.evicted_records)
.map(|outcome| ByteSize::b(outcome.evicted_bytes as u64))

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.

ByteSize for the win!

@nadav-govari
nadav-govari force-pushed the feature/rework-scaler branch from cca58bd to a906e6e Compare October 5, 2026 20:57

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants