Skip to content

perf(internal_logs source): decouple broadcast drain from downstream send and batch events - #26518

Open
thomasqueirozb wants to merge 3 commits into
internal-logs-dropfrom
internal-logs-throughput
Open

thomasqueirozb wants to merge 3 commits into
internal-logs-dropfrom
internal-logs-throughput

Conversation

@thomasqueirozb

@thomasqueirozb thomasqueirozb commented Sep 30, 2026 •

Copy link
Copy Markdown
Member

Summary

Stacked on #25218, which makes internal_logs broadcast lag visible through component_discarded_events_total{intentional="false"}. This PR reduces how often that lag occurs.

  • Decouples broadcast consumption from downstream sending. A dedicated drain task pulls from the trace broadcast into a bounded intermediate queue. The main task batches from the queue with recv_many and calls send_batch. This keeps the broadcast receiver drained while the sink is backpressured, and amortizes per-event overhead downstream.
  • The intermediate queue holds one batch (MAX_BATCH_SIZE, 1024 events). A 10,000-event queue lowered the console drop rate by only 0.2 percentage points (benchmark) and raises worst-case memory 10x.
  • The main task keeps a clone of the ShutdownSignal for the full run() scope, so the source does not report completion while batches are still in flight.
  • If send_batch fails, the drain task is aborted before the source returns, so it cannot keep its ShutdownSignal clone alive.

Moved from #25218. The review threads there about INTERMEDIATE_QUEUE_CAPACITY and drain task cleanup when send_batch fails are addressed in 5604fed.

Vector configuration

Minimal config (from the issue, console sink):

api:
  enabled: true

sources:
  internal_logs:
    type: internal_logs
  internal_metrics:
    type: internal_metrics
    scrape_interval_secs: 1

sinks:
  show_internal_logs:
    type: console
    inputs:
      - internal_logs
    encoding:
      codec: json
  prom:
    type: prometheus_exporter
    inputs:
      - internal_metrics
    address: 127.0.0.1:9598

Blackhole sink (isolates the source path from stdout/JSON costs):

api:
  enabled: true

sources:
  internal_logs:
    type: internal_logs
  internal_metrics:
    type: internal_metrics
    scrape_interval_secs: 1

sinks:
  null_sink:
    type: blackhole
    inputs:
      - internal_logs
  prom:
    type: prometheus_exporter
    inputs:
      - internal_metrics
    address: 127.0.0.1:9598

Benchmark

The prometheus_exporter is scraped at the end of each 20s run to read component_received_events_total and component_discarded_events_total{intentional="false"} for the internal_logs source.

"Single-task loop" and "master (patched)" rows use master's source loop with lag counted in the drop metric, which is the behavior after #25218.

The tables below were measured with an intermediate queue capacity of 10,000. The final capacity of 1024 changed the console drop rate by 0.2 percentage points in a separate run (details).

Design comparison, console sink, VECTOR_LOG=trace, 20s

Buffer size is the broadcast capacity in src/trace.rs.

Design Broadcast buffer Drops
Single-task loop (master) 99 876,567
Single-task loop 10,000 848,510
Drain + batching 99 353,217
Drain + batching 10,000 333,105

Buffer size made almost no difference (3-6% fewer drops) with either design, so the original 99 is retained.

Sink comparison, VECTOR_LOG=trace, 20s

Version Sink Received Dropped Total Drop %
master (patched) console 129,468 776,664 906,132 85.7%
master (patched) blackhole 147,279 883,563 1,030,842 85.7%
this branch console 395,443 353,870 749,313 47.2%
this branch blackhole 1,524,445 0 1,524,445 0%

Interpretation:

  • With the single-task loop the source itself is the bottleneck: single-event send_event and broadcast consumption being coupled cap throughput at ~5k events/sec delivered and drop ~86% of events even when the sink is free (blackhole).
  • With the drain + batching design, the source can deliver ~76k events/sec (~10x higher delivered throughput, ~1.5x higher combined throughput) when the sink doesn't backpressure. On the console sink it still drops under trace because stdout + JSON encoding caps at ~20k events/sec.

Under VECTOR_LOG=debug (normal load), both configs show zero drops.

How did you test this PR?

  • cargo nextest run --no-default-features --features sources-internal_logs --lib sources::internal_logs:: (all tests pass, including broadcast_lag_increments_discarded_metric from fix(internal_logs source): report broadcast lag drops in component_discarded_events_total #25218)
  • cargo clippy --no-default-features --features sources-internal_logs --lib --tests -- -D warnings
  • cargo vdev check events
  • Ran both configs at VECTOR_LOG=debug and VECTOR_LOG=trace, comparing component_received_events_total and component_discarded_events_total (see Benchmark).

Change Type

  • Bug fix
  • New feature
  • Dependencies
  • Non-functional (chore, refactoring, docs)
  • Performance

Is this a breaking change?

  • Yes
  • No

Does this PR include user facing changes?

  • Yes. Please add a changelog fragment based on our guidelines.
  • No. A maintainer will apply the no-changelog label to this PR.

References

@thomasqueirozb
thomasqueirozb added this pull request to stack #26519 September 30, 2026 19:14
@thomasqueirozb thomasqueirozb changed the title enhancement(internal_logs source): decouple broadcast drain from downstream send and batch events perf(internal_logs source): decouple broadcast drain from downstream send and batch events Sep 30, 2026
@thomasqueirozb
thomasqueirozb marked this pull request as ready for review September 30, 2026 19:15
@thomasqueirozb
thomasqueirozb requested a review from a team as a code owner September 30, 2026 19:15

@datadoghq-integration datadoghq-integration 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.

Bits Code Review: PASS

More details

The bounded drain task, batching loop, shutdown guard, and abort-on-send-failure lifecycle are internally consistent; no changed-line failure mode was identified.

Was this helpful? React 👍 or 👎

Open Bits AI session

🤖 Bits Code Review · Commit b535353 · @DataDog review to ask questions

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

domain: sources Anything related to the Vector's sources

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant