Skip to content

Avoid broadcast joins based on incomplete empty samples - #24260

Open
pentschev wants to merge 1 commit into
NVIDIA:mainfrom
pentschev:cudf-polars/incomplete-zero-sample-guard
Open

pentschev wants to merge 1 commit into
NVIDIA:mainfrom
pentschev:cudf-polars/incomplete-zero-sample-guard

Conversation

@pentschev

Copy link
Copy Markdown
Contributor

Description

Dynamic join planning estimates input sizes from an initial sample of chunks. If every sampled chunk contains zero rows, the current logic treats that input as empty and eligible for broadcasting.

An all-zero sample is conclusive only when sampling consumed the complete input. For an incomplete sample, later chunks may contain substantial data, particularly when reading ordered files after predicate pushdown.

This change treats an incomplete all-zero sample as having an unknown size for broadcast eligibility. An estimate remains eligible when:

  • sampling is complete; or
  • the sampled prefix contains nonzero rows or bytes.

Complete empty inputs and incomplete nonzero samples retain their existing behavior.

Motivation

This issue was observed with a TPC-H dataset containing many small, ordered Parquet files.

For SF300 Q3, the first 41 of 90 relevant lineitem chunks contain no rows after applying the l_shipdate > 1995-03-15 predicate. Sampling only the initial prefix therefore produced an incomplete zero-size estimate.

Before this change, dynamic planning selected the large lineitem input as the broadcast side:

  • Decision: broadcast_right
  • Join evaluations: 2,639
  • Cumulative evaluated input: approximately 1.08 TiB
  • Traced runtime: 9.2330 seconds

With this change, the incomplete zero-size estimate is not considered a valid broadcast candidate:

  • Decision: broadcast_left
  • Join evaluations: 91
  • Cumulative evaluated input: approximately 95.4 GiB
  • Traced runtime: 4.9317 seconds

The focused regression test reproduces the same strategy-selection difference:

  • Before: BroadcastJoinStrategy(side="right")
  • After: BroadcastJoinStrategy(side="left")

End-to-end results

The isolated comparison used #24209 at commit
737734dcedc85d8bb8cfe26d7eb67dc1bb937648 with the following additional settings:

CUDF_POLARS__PARQUET_OPTIONS__MAX_FOOTER_SAMPLES=3
CUDF_POLARS__EXECUTOR__DYNAMIC_PLANNING__SAMPLE_CHUNK_COUNT=8

PR #24209 without this change

  • SF300:
    • 22/22 queries completed
    • 56.2321-second query sum
    • 83,648 MiB GPU peak
  • SF1000:
    • Only 20/22 queries completed
    • Q3 and Q7 failed with CUDA allocation OOMs, 147.4632-second query sum on queries that didn't fail
    • 96,980 MiB GPU peak

PR #24209 with this change

  • SF300:
    • 22/22 queries completed
    • 56.0651-second query sum
    • 78,430 MiB GPU peak
  • SF1000:
    • 22/22 queries completed
    • 174.9092-second query sum
    • 81,438 MiB GPU peak

Tradeoffs

This makes strategy selection more conservative when an incomplete sampled prefix contains no data.

If the remaining unsampled input is nonempty but genuinely small, the planner may decline to broadcast that side and instead broadcast the other side or select a shuffle join. Such workloads could therefore perform more work than before.

The affected sample currently provides no upper bound for the unsampled data, however, so treating it as broadcastable can result in arbitrarily underestimating its size. This change favors a bounded decision over that potentially unsafe assumption.

@pentschev pentschev self-assigned this Sep 22, 2026
@pentschev pentschev added the 3 - Ready for Review Ready for review by team label Sep 22, 2026
@pentschev
pentschev requested a review from a team as a code owner September 22, 2026 18:30
@pentschev pentschev added the improvement Improvement / enhancement to an existing function label Sep 22, 2026
@pentschev pentschev added non-breaking Non-breaking change cudf-polars Issues specific to cudf-polars labels Sep 22, 2026
@github-actions github-actions Bot added the Python Affects Python cuDF API. label Sep 22, 2026
@coderabbitai

coderabbitai Bot commented Sep 22, 2026

Copy link
Copy Markdown

Review in Change Stack →

Navigate logical layers of code changes, visualize relationships, and explore their blast radius.

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Repository: NVIDIA/cudf/.coderabbit.yaml

Review profile: CHILL

Plan: Enterprise

Run ID: 31e473b9-315a-475a-9c77-b830a5e929e1

📥 Commits

Reviewing files that changed from the base of the PR and between b11327a and 5838fdc.

📒 Files selected for processing (2)
  • python/cudf_polars/cudf_polars/streaming/actor_graph/join.py
  • python/cudf_polars/tests/streaming/test_join.py

Included review availability: Your plan provides up to 12 included reviews per hour; 11 remain after this review.


📝 Summary

Summary by CodeRabbit

  • Bug Fixes
    • Improved dynamic join strategy selection when sampled input statistics are incomplete or contain zero rows.
    • Prevented incomplete, all-zero samples from being incorrectly treated as empty inputs eligible for broadcasting.
    • Preserved existing byte, row-count, and duplication limits for broadcast decisions.

Walkthrough

Changes

Join broadcast selection

Layer / File(s) Summary
Sample eligibility validation
python/cudf_polars/cudf_polars/streaming/actor_graph/join.py
_choose_strategy_from_samples now rejects incomplete samples with zero bytes and rows. Existing byte, row-count, and duplication limits remain unchanged.
Regression coverage
python/cudf_polars/tests/streaming/test_join.py
The test setup imports the required strategy and metadata types. A regression test verifies that an incomplete zero-row right sample results in left-side broadcasting.

Priority: ➖ Normal

Estimated code review effort: 2 (Simple) | ~10 minutes

Change: Bug fix

Merge Risk: ⚪ Minimal · up to 5838f

No actionable merge-blocking risk is identified from the reviewed change.

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 60.00% which is insufficient. The required threshold is 80.00%. Docstring coverage is scoped to functions touched by this diff. Analyzed 5 functions across 2 files. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Title check ✅ Passed The title clearly and concisely describes the main change: preventing broadcast joins when samples are incomplete and empty.
Description check ✅ Passed The description directly explains the incomplete-sample issue, the eligibility rules, the regression test, motivation, results, and tradeoffs.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
  • Fix all pre-merge checks with AI
✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create a new PR

Comment @coderabbitai help to get the list of available commands.

Comment on lines +1238 to 1254
# An incomplete all-zero prefix is not evidence that the full input is empty.
# Treat it as unknown rather than choosing it as the broadcast side. This can
# happen when an ordered scan's early chunks are eliminated by a predicate.
left_size_known = left_sample.is_complete or left_total > 0 or left_total_rows > 0
right_size_known = (
right_sample.is_complete or right_total > 0 or right_total_rows > 0
)
right_size_ok = right_total < broadcast_threshold and (
right_total_rows < MAX_ROWS_PER_PARTITION or right_metadata.duplicated
left_size_ok = (
left_size_known
and left_total < broadcast_threshold
and (left_total_rows < MAX_ROWS_PER_PARTITION or left_metadata.duplicated)
)
right_size_ok = (
right_size_known
and right_total < broadcast_threshold
and (right_total_rows < MAX_ROWS_PER_PARTITION or right_metadata.duplicated)
)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

This feels like a sticking plaster. We don't like to consume too much of an input channel because we don't want to buffer unboundedly. But if the messages are empty frames then that isn't actually a memory problem.

Additionally, sending loads of empty messages is going to reduce performance because we have to schedule and execute this work that produces no output.

So can we somehow address the root cause of this problem? From what I understand, if the input table is ordered then the filter that we apply can produce zero rows in the scan and we send that downstream. Because we basically don't repartition at any point we have this long prefix of empty tables before we get to the part of the scanned channel that has any data in it.

It seems like we should never have sent those empty messages on at all.

We currently rely on channels always sending at least one message (even if it is empty) so that the consumer gets the right schema-d table in various places. But we probably shouldn't send loads.

I think there are only a few places where this can happen, so maybe we can do this differently...

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Hmm, I thought a bit harder and I think squeezing empty messages out of a channel is a no-go for now without reworking the way we match up sides of a shuffled join. Today we require that if we're doing a shuffled join that the sequence numbers of the two sides of the input match up (and we pair them in order). If we were to squeeze out empty messages we would need a different contract for matching message where a gap in sequence numbers is indicative of an empty chunk (which means in a shuffle join the other side would not match).

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Another thing we could do though is to remove the "max sample chunks" from the prefix sampler and sample up to some byte count.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Did that here: #24270

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

Labels

3 - Ready for Review Ready for review by team cudf-polars Issues specific to cudf-polars improvement Improvement / enhancement to an existing function non-breaking Non-breaking change Python Affects Python cuDF API.

Projects

Status: Todo

Development

Successfully merging this pull request may close these issues.

2 participants