Skip to content

Navigation Menu

Sign in
Sign up

[Ideas] Implement Robust Predicate Transfer (RPT+) on top of Anser #1963

Unanswered
leborchuk asked this question in Ideas / Feature Requests
Discussion options

Description

Implement Robust Predicate Transfer with dynamic execution (RPT+) on top of
Anser (gpcontrib/anser/). Read gpcontrib/anser/README.md for the current
architecture before starting.

The references, precisely

What Where Note
The paper: RPT+ Yiming Qiao, Peter Boncz, Huanchen Zhang, Robust Predicate Transfer with Dynamic Execution, PVLDB 19(6):1278-1290, 2026, doi:10.14778/3797919.3797934 read this one
The paper's own artifact https://github.com/embryo-labs/dynamic-predicate-transfer this is RPT+
The original RPT extension https://github.com/robust-sql/robust this is the baseline RPT+ improves on, not RPT+

Read both repos, but know which is which: robust-sql/robust implements plain
Predicate Transfer (DuckDB's native bloom filter, 12 bits/key; reports ×ばつ
geomean on JOB), while the paper is specifically about fixing RPT's overhead
and regressions
.

⚠️ Licensing. Cloudberry is Apache-2.0 and AGENTS.md in this repo says:
"Do not generate or import code with incompatible licensing." Check the license
of both repos before reusing even a snippet, and prefer implementing from the
paper. If you do reuse anything, it needs attribution plus LICENSE/NOTICE
handling — raise it in the PR rather than deciding alone.

What the paper actually proposes

RPT (the baseline) makes Yannakakis' algorithm practical: model the query as
a join graph G(V,E) where edge weight = number of shared join attributes,
build a maximum spanning tree rooted at the largest table (LargestRoot),
then propagate bloom filters over that tree in a forward pass (leaves→root)
and a backward pass (root→leaves), fully reducing tables before the join
phase. Its weakness: it ignores equivalence relationships between join keys,
producing redundant and oversized filters, and pays heavy scan +
BF-construction overhead — enough to cause real regressions.

RPT+ adds three things:

  1. Asymmetric Transfer Plan (ATP) — different tree shapes for the two
    passes, because they do different jobs. Forward = collect information, so a
    deep chaining tree is best (filter table i before building the filter
    for i+1, avoiding oversized BFs). Backward = distribute information, so a
    wide broadcast tree is best (build the final filter once and reuse it
    for every table in the same equivalence class). Chains for different
    equivalence classes are linked through "bridge tables" (the root of each
    class's induced subtree).
  2. Cascade filters — min/max for block-level skipping plus a bloom filter
    for tuple-level filtering, applied coarse-to-fine.
  3. Dynamic pipelines — runtime selectivity monitoring that abandons filter
    construction when it will not pay off (Algorithm 2 in §5.2).

Reported speedups over DuckDB v1.3.0: ×ばつ JOB, ×ばつ SQLStorm, ×ばつ TPC-H,
×ばつ Appian
, "avoid[ing] the significant performance regressions observed
with the original RPT".

Constants from §5.5, use these as starting values, do not invent your own:

Parameter Value Meaning
γ 100 K tuples rows observed before deciding
τ_sel 0.35 abandon if > 35% of tuples survive
τ_prog 0.6 only abandon while < 60% through the scan
M_avail 64 GB memory budget for materialization
τ_stop 0.9 stop probing a filter that rejects too little
filter 20 bits/key, 7 hash functions FPR ×ばつ10−5
probe cost 2.48 cycles/tuple 64-byte cache-line blocks, 32-bit sectors, VPGATHERDD

Sampling is sequential, not random — random access breaks streaming
execution. On JOB, τ_sel in 0.1-0.3 gave the best speedups; τ_sel = 0 (drop
almost everything) hurt.

Why Anser fits — and the one place it does not

RPT+ concept Anser today Gap
equivalence class of join attributes condition_keyinclude/anser.h documents it as exactly "the optimizer-generated equivalence-class symbol" naming only
backward pass: build the final filter once, broadcast to every table in the class one channel per condition, merged once in the QD backend, pushed to every subscriber (anser_disp_deliver(), src/anserdispatch.c:372) already implemented
combining filters built by different workers AnserBloomFoldPartInPlace() OR-folds parts as they arrive (src/anserfilter.c:126) already implemented (the paper is single-node; this is our addition)
forward pass: chaining, table i reduced before it builds the filter for i+1 channels are independent, with no ordering between them the hard part — see below
dynamic pipeline abandonment (Alg. 2) never abandons overlaps the bloom-performance issue, task 2.4
cascade filter, blocked/sectorized BF classic bloom: k random bitset accesses overlaps the bloom-performance issue, tasks 2.1/2.3
one filter per (table, equivalence class) plan pass injects one producer+consumer per hash join needs join-graph analysis in src/anserplan.c

Anser's backward pass is essentially free — it is what the extension already
does, and the dispatch-connection transport made it cheaper still: one merge in
the coordinator backend, then one push per consumer, no shared memory and no
worker in between. The forward chaining pass is where this project lives or
dies
, for a reason that does not exist in DuckDB:

  • In DuckDB, a chain hop is an in-process pipeline breaker: microseconds of
    coordination.
  • In Anser, every hop is a global barrier: all segments publish → coordinator
    merges → all segments receive.

And we now know what a hop costs. A traced exchange on a 3-segment demo cluster
with a 1 MB filter took 34 ms, of which ~21 ms was the coordinator picking
parts up one at a time (~7.2 ms per part, gated by the interconnect wait loop).
JOB queries join 5-10 tables, i.e. 4-9 sequential hops — on the order of
140-300 ms of pure barrier latency added to queries whose total runtime
may be under a second, on loopback, with the smallest filter Anser can build.
That arithmetic is the single most important input to milestone 0.

⚠️ Correctness rule you must not get wrong. It is tempting to do the forward
chaining locally on each segment and avoid the barrier. That is only correct
when the tables are co-distributed on the join key.
Otherwise segment s
holds only part of table j's key set, and filtering table j+1 with it drops
rows whose key lives on another segment — wrong results, silently. Where
tables are not co-located, each hop must be global. Write this rule into the
design note, with the co-location test, before any code.

Milestones

M0 — Read, map, measure the budget (blocking; no implementation)

Design-note deliverable:

  1. A corrected version of the mapping table above, after reading the paper and
    both repos.
  2. Measured single-hop Anser latency (publish → merge → deliver) at 3, 8 and
    24 segments. anser.debug already logs every step with timestamps, so this
    is a measurement, not new instrumentation. Then the projected cost of a 4-hop
    and 8-hop chain. State whether M3 is viable.
  3. The co-location rule: how to detect from the plan that two tables are
    distributed on the same key, and how many JOB / TPC-DS joins qualify.
  4. Where the join graph, equivalence classes and predicates come from. Anser's
    pass runs post-plan (AnserApplyRuntimeFilters, src/anserplan.c), which
    sees a finished PlannedStmt, not the join graph — can equivalence classes be
    recovered there, or does this need a pre-plan hook? This is the second
    make-or-break question.
    ORCA and the Postgres planner will differ; report
    both.
  5. Memory, not channel count. Channels are a hash in the coordinator
    backend now, so there is no anser.max_channels to size — but RPT creates one
    filter per (table, class), so a 10-table query could hold ten 1 MB
    accumulators per query in the QD, plus libpq's per-connection input buffers
    (which keep their high-water mark for the session). Budget that, and say what
    the cap should be.

Do not proceed past M0 without review. "The forward pass costs more than it
saves in an MPP" is a legitimate and valuable outcome.

M1 — Backward broadcast pass across an equivalence class (start here)

Generalize today's one-producer-set/one-consumer-per-hash-join to: all tables
in an equivalence class publish, the merge happens once, every table in that
class receives it.
No ordering, no barrier beyond the one Anser already has,
and it is exactly the pass RPT+ shows is best served by a broadcast tree. Expect
the bulk of the achievable win here.

M2 — Dynamic pipeline (Algorithm 2)

Implement the abandon logic with the paper's constants, on the producer side.
Note this is the same mechanism as task 2.4 of the bloom-performance issue
coordinate, do not implement it twice; whoever gets there first owns it, the
other reviews. Anser already has the cancel path
(ExecAnserBloomFilterProduceCancel(), src/anserbloomproduce.c), but mind the
semantics: an Anser cancel kills the channel for every consumer, whereas RPT+
abandons one filter locally. Decide and document which you want.

M3 — Asymmetric transfer plan (forward chaining) — only if M0 allows

Chaining forward tree + broadcast backward tree, bridge tables between classes.
Restrict chain depth by what M0's latency budget supports, and say what depth
you allowed and why.

M4 — Cascade filter

min/max block skipping + a blocked/sectorized bloom layout. Cross-link to the
bloom-performance issue: RPT+ affords 20 bits/key and 7 hashes precisely
because a blocked filter touches one cache line per probe (2.48
cycles/tuple), while Anser's classic bloom does k random accesses. That suggests
the layout matters more than the bits/key knob — and a blocked filter is still
unionable by bitwise OR when both sides share parameters, so the cheap fold
survives. Verify that claim before relying on it.

Benchmarks

  • JOB (Join Order Benchmark, IMDb) is mandatory — it is the paper's headline
    workload and this technique's home turf. Loading IMDb into Cloudberry with
    sensible distribution keys is part of the task; commit the DDL under
    gpcontrib/anser/doc/bench/.
  • TPC-DS SF100 minimum, using the reporting protocol from the
    bloom-performance issue (injected subset in detail, non-injected regression
    check, geomean).
  • Report regressions first, wins second. RPT's entire failure mode was
    regressions, and RPT+'s main claim is removing them; a report that shows the
    geomean win without a per-query regression table does not answer the question
    this work exists to answer.
  • Compare against the paper's numbers (×ばつ JOB) as a sanity check, but expect
    to differ: we are distributed, our hops cost a global barrier, and our filter
    implementation is not theirs. Explain your gap rather than matching it.
  • Setup: shared_preload_libraries='anser', anser.enable=on,
    anser.runtime_filter=on. No CREATE EXTENSION needed — the subsystem
    creates no catalog objects.

Definition of done

  • M0 design note, reviewed: mapping, measured hop latency and chain budget,
    co-location correctness rule, join-graph feasibility for both optimizers,
    per-query memory budget.
  • M1 implemented behind a GUC, default off, fail-open on every failure path,
    visible in EXPLAIN ANALYZE.
  • Results identical to the unmodified plan: full result-set comparison (not
    row counts) on the regression suite, JOB and TPC-DS.
  • JOB + TPC-DS reports in gpcontrib/anser/doc/performance.md, regressions
    reported per query.
  • Licensing question answered explicitly in the PR.

Out of scope

  • Cyclic queries. Yannakakis, RPT and RPT+ all target acyclic (α-acyclic)
    queries. Detect and skip the rest; do not try to generalize.
  • The full semi-join reduction guarantee — RPT+ itself trades it away
    ("restricts the forward transfer plan to tables with predicates, plus bridge
    tables ... trades the theoretical guarantee for significantly reduced
    overhead"). We are following it there, not doing better.
  • Rewriting the join order or the optimizer. This technique exists to make a
    bad join order survivable, not to fix it.
  • Porting DuckDB code. See the licensing note.

Use case/motivation

No response

Related issues

#1942

Are you willing to submit a PR?

  • Yes I am willing to submit a PR!
You must be logged in to vote

Replies: 0 comments

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment
Labels
type: Enhancement New feature or request, ideas good first issue Good for newcomers
1 participant

AltStyle によって変換されたページ (->オリジナル) /