Hi Asif, Thanks for the careful read -- these are exactly the right questions.
> would it mean n calls to fetch the hashed partitions? Or ... identify > the needed partitions in a single call, as multiple hash partitions > may reside on a single executor? You're right that we batch, but let me be precise about the granularity, because it's per-shard, not per-executor. On a probe task, each incoming row is hashed to its target shard and buffered into a per-shard buffer. When a shard's buffer reaches maxBatchSize (default 1024 keys) it's flushed as one async RPC to the executor holding that shard; up to maxInFlightNum (default 8) RPCs are in flight concurrently per task, and any partially-filled buffers are flushed at the end. So a single probe row is never an individual RPC -- it's amortized into a batch. But the batch is scoped to one shard. So if a probe task has keys spanning all N shards, it issues at least N batched RPCs (more when a shard's key count exceeds maxBatchSize). We do NOT currently coalesce multiple shards that happen to be co-located on the same executor into a single call. That's a valid optimization -- group batches by target executor rather than by shard -- and a good candidate for a follow-up; it just isn't in the current implementation. Good catch. > In the worst case scenario of all partitioned keys are required by all > the probing nodes, it degenerates to BHJ, right? Not quite -- and the difference is worth spelling out. BHJ replicates the entire build side to every executor (one full local copy each) and does zero network I/O at probe time. DMJ never puts a full copy on the probe side; it fetches matched build rows on demand and streams them back. So in the pathological case you describe -- Bloom filter gives no benefit and every probe key matches -- DMJ does not reduce to BHJ. Two things happen instead: (1) the RPC topology becomes all-to-all (every probe executor talks to every shard executor), which looks more like a shuffle than a broadcast; and (2) a hot build row that matches K probe rows is transmitted once per match, so total bytes moved can actually exceed BHJ. In other words, that's the case where DMJ is the *wrong* choice -- if the build side fits, BHJ wins; otherwise a shuffle join wins. That's precisely why DMJ is hint-only with no cost-based selection: the planner never picks it for you. It's designed for the opposite regime -- a medium build side (~200 MB-10 GB) with a very large, selective probe side, where the Bloom filter drops 60-90% of probe rows before any RPC and the win comes from not shuffling the huge probe side. Push it toward the all-keys-match extreme and it stops being the right tool. Thanks again -- happy to go deeper on either point. Yu Gan On Thu, Aug 6, 2026 at 5:14 PM Asif Shahid <[email protected]> wrote: > > Also another clarification requested > In the worst case scenario of all partioned keys are required by all the > probing nodes, it degenerates to BHJ, right? > > democracy of barbarians is worse than dictatorship > > On Thu, Aug 6, 2026, 2:04 AM Asif Shahid <[email protected]> wrote: >> >> Just went through the doc.. >> Based on that it seems that you will be batching to minimize the calls. >> Pls correct me if wrong >> >> democracy of barbarians is worse than dictatorship >> >> On Thu, Aug 6, 2026, 1:57 AM Asif Shahid <[email protected]> wrote: >>> >>> Interesting. >>> One question: >>> Assume that there are say n hash partitions of build side data. >>> On a given node the streaming side( probe side ) has keys corresponding to >>> to those n partitions.. would it mean n calls to fetch the hashed >>> partitions? >>> >>> Or you plan to group the probe side in batches, identify the needed >>> partitions in a single call >>> As multiple hash partitions may be residing on a single executor. >>> Regards >>> Asif >>> >>> democracy of barbarians is worse than dictatorship >>> >>> On Thu, Aug 6, 2026, 1:44 AM Yu Gan <[email protected]> wrote: >>>> >>>> Hi all, >>>> >>>> I'd like to start a discussion on a SPIP for a new opt-in join >>>> strategy: Distributed Map Join (DMJ). >>>> >>>> - Design doc (SPIP format, open for comments): >>>> >>>> https://docs.google.com/document/d/1ioj2vaY1-kgvGX-FCCVA95ybBYjlnc20KePS1GNPbSM >>>> - JIRA: https://issues.apache.org/jira/browse/SPARK-57487 >>>> - PR (reference implementation): https://github.com/apache/spark/pull/56542 >>>> - Shepherd: Chao Sun (sunchao) >>>> >>>> Problem. There's a gap between broadcast join and shuffle join. When >>>> the build side is medium-sized -- too large to broadcast (roughly >>>> hundreds of MB up to ~10 GB) -- but the probe side is very large >>>> (PB-scale in our production case), Spark falls back to a shuffle join, >>>> and shuffling the probe side becomes the dominant cost of the query. >>>> >>>> Proposal. DMJ avoids the probe-side shuffle. The build side is >>>> hash-partitioned into N remotely-queryable shards held in-memory on >>>> executors; probe tasks send batched RPC lookups to the shard-holding >>>> executors, pre-filtered by a Bloom filter, and stream matched rows >>>> back. No probe-side shuffle. >>>> >>>> Scope / safety. The feature is double-gated and off by default: >>>> - spark.shard.enabled=true starts the shard infrastructure at >>>> application launch (inlined in the executor process -- an additional >>>> Netty server + RPC endpoints in SparkEnv, no separate process or >>>> deployment step). >>>> - An explicit per-query SQL hint /*+ DISTMAPJOIN(table) */ activates >>>> the strategy. >>>> >>>> There is no cost-based selection; the planner never chooses DMJ >>>> automatically. When disabled, no code paths are initialized and there >>>> is zero overhead. >>>> >>>> Production status. Running stably in our 16k-core production >>>> environment for 2+ months on PB-scale workloads (~60% wall-time >>>> reduction on qualifying fact-dimension joins). >>>> >>>> The design doc covers architecture, data flow, the key design >>>> decisions, risks, and rejected alternatives. I'd appreciate feedback >>>> on the overall approach and on the specific design questions raised on >>>> the PR (e.g. whether the shard infrastructure should remain opt-in, >>>> durability/fault-tolerance model, and transport/auth). >>>> >>>> Thanks, >>>> Yu Gan >>>> >>>> --------------------------------------------------------------------- >>>> To unsubscribe e-mail: [email protected] >>>> --------------------------------------------------------------------- To unsubscribe e-mail: [email protected]
