Skip to content

Reduce GIL usage - #376

Open
fpacifici wants to merge 4 commits into
mainfrom
worktree-reduce_gil
Open

fpacifici wants to merge 4 commits into
mainfrom
worktree-reduce_gil

Conversation

@fpacifici

Copy link
Copy Markdown
Contributor

The rust consumer moves the payload of each message into python memory immediately after getting the message from Kafka.

This was ok when all messages are processed by Python. Taking the GIL and moving memory was irrelevant.
Now we do a number of processing steps directly in rust like header filters and batching.
Though we did not remove the python processing that happens before filtering.
In prod we observed consumers underperforming when they filter out most messages.
Taking the GIL for each message that is filtered out could be the reason.

This PR introduces Raw rust messages processable by rust native steps without moving data to
Python first.

Filippo Pacifici and others added 4 commits September 13, 2026 16:39
Five steps destructured the payload with `let ... else` or `if let`, so a new
enum variant would compile without error and silently take the else branch:
the header filter would stop filtering, map would become a no-op, routing
would be skipped, watermarks would stop tracking the last message time, and
the Python filter would panic on the first message.

Spell out every arm instead, so adding a variant is a build failure. No
behaviour change.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
`RawMessage.payload` becomes `Arc<[u8]>` so cloning a message stays cheap:
`Broadcaster` clones once per downstream branch, and an upcoming Rust-native
payload variant will be cloned along the whole chain. It also makes the
immutability this module already documents real.

Alongside it, remove code that has no callers:
- `#[pyo3(set)]` on `RawMessage.timestamp`/`.schema`. Nothing assigns them from
  Rust or Python, and rust_streams.pyi already declares both read-only.
- the unused `StreamingMessage` enum.
- `unwrap_payload()` is narrowed to `#[cfg(test)]`; production code matches on
  every variant instead.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
`RoutedValuePayload::RustRawMessage` carries a `RawMessage` that lives in Rust
memory, so a step that can do its work natively never takes the Gil. Every
match on the enum gets an arm for it:

- watermark and the header filter read the timestamp and the headers natively.
  Both run before any filtering, on 100% of the traffic; the header filter also
  stops cloning the header vector per message.
- batch stores the variant as it arrives (`Vec<RoutedValuePayload>`, mixed) and
  only builds the Python list on flush.
- sinks and the GCS writer read the bytes natively instead of copying them out
  of a PyBytes.
- map and PythonAdapter convert and ratchet: they already replace the payload,
  so the Python form is what continues downstream.
- filter and router convert transiently and forward the original Rust form, so
  a following Rust step keeps the cheap path.

Nothing emits the variant yet; flipping the source is the next commit.

Test builders default to the Rust representation, so the existing suite
exercises what production will run, with an explicit Python-representation test
for each step that branches on it.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
`to_routed_value` no longer allocates a `Py<RawMessage>` per message. The payload
stays in Rust memory until a step hands it to Python code.

On a pipeline shaped like header_filter -> batch -> Python work, where the
header filter drops most events, the source, the watermark step and the filter
now run without taking the Gil at all. A message the filter drops never enters
Python memory, so it costs zero Gil acquisitions instead of three.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@fpacifici
fpacifici requested a review from a team as a code owner September 14, 2026 00:38
Comment on lines +524 to +526
fn from(value: &RawMessage) -> Self {
PyStreamingMessage::RawMessage {
content: traced_with_gil!(|py| into_pyraw(py, value.clone()).unwrap()),

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Bug: The use of .unwrap() when converting a RawMessage to a Python object can cause a panic if the object allocation fails, for example due to an out-of-memory error.
Severity: MEDIUM

Suggested Fix

Replace the .unwrap() calls with proper error handling. The From implementations should propagate the PyResult from into_pyraw instead of panicking. This will allow the calling context to handle potential allocation failures gracefully, for instance by returning an error that can be managed further up the call stack.

Prompt for AI Agent
Review the code at the location below. A potential bug has been identified by an AI
agent. Verify if this is a real issue. If it is, propose a fix; if not, explain why it's
not valid.

Location: sentry_streams/src/messages.rs#L524-L526

Potential issue: The `From<&RawMessage>` implementations for `Py<PyAny>` and
`PyStreamingMessage` use `.unwrap()` on the result of the `into_pyraw` function. This
function can fail during Python object creation, for example, in an out-of-memory (OOM)
scenario. Because the conversion from `RawMessage` to a Python object is a common
operation in production pipelines (e.g., in map, filter, and router steps), a failure in
`Py::new()` will result in an unhandled panic, causing the process to crash instead of
handling the error gracefully.

Did we get this right? 👍 / 👎 to inform future reviews.

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