恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
Daft Join 策略完全指南:从七种 Join 类型到 Hash / Broadcast / Sort-Merge 执行策略与分布式调优
首页
资讯中心
/
Daft Join 策略完全指南:从七种 Join 类型到 Hash / Broadcast / Sort-Merge 执行策略与分布式调优
Daft Join 策略完全指南:从七种 Join 类型到 Hash / Broadcast / Sort-Merge 执行策略与分布式调优
发布时间:2026/9/18 3:30:57
Daft Join 策略完全指南从七种 Join 类型到 Hash / Broadcast / Sort-Merge 执行策略与分布式调优【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/DaftDaft 的DataFrame.join()支持七种 Join 类型inner / left / right / outer / semi / anti / cross与三种执行策略hash / broadcast / sort_merge。Join 类型决定输出中包含哪些行执行策略决定 Daft 在物理上如何执行 Join——哪一侧被 Shuffle、被广播或被排序。本文以 官方 Join 策略文档 为核心骨架结合仓库源码daft/dataframe/dataframe.py、src/daft-core/src/join.rs、src/daft-logical-plan/src/optimization/rules/reorder_joins/等逐层展开帮助你理解每种组合的适用场景并掌握在单机与 Ray 分布式环境下的内存调优手段避免次秒级操作与 OOM 崩溃之间的天壤之别。一、七种 Join 类型选择输出行的规则Join 类型的语义与 SQL 完全对齐。在 join.rs 中JoinType枚举定义了六种核心类型pub enum JoinType { Inner, Left, Right, Outer, Anti, Semi, }而 Python 层的how参数支持的第七种cross在 字符串解析 中被映射为JoinType::Inner——正如源码注释所说cross join is just inner join with no join keysCross Join 就是不带 Join Key 的 Inner Join。下表汇总了全部七种类型及其 SQL 对照TypeSQL EquivalentOutputinnerINNER JOINRows where both sides match on the join keyleftLEFT OUTER JOINAll rows from the left, with nulls where the right has no matchrightRIGHT OUTER JOINAll rows from the right, with nulls where the left has no matchouterFULL OUTER JOINAll rows from both sides, with nulls where either side has no matchsemiWHERE EXISTSRows from the left that have at least one match on the rightantiWHERE NOT EXISTSRows from the left that have no match on the rightcrossCROSS JOINCartesian product of both sides (no join key)join()的完整签名daft/dataframe/dataframe.py#L3921-L3931如下其中on/left_on/right_on用于指定连接键def join( self, other: DataFrame, on: list[ColumnInputType] | ColumnInputType | None None, left_on: list[ColumnInputType] | ColumnInputType | None None, right_on: list[ColumnInputType] | ColumnInputType | None None, how: Literal[inner, left, right, outer, anti, semi, cross] inner, strategy: Literal[hash, sort_merge, broadcast] | None None, prefix: str | None None, suffix: str | None None, ) - DataFrame1.1 连接键的三种指定方式on左右两侧键名相同最常见left_onright_on两侧键名不同时分别指定cross不允许传入任何键。源码在 dataframe.py#L3978-L3992 中做了严格校验crossJoin 中设置on/left_on/right_on或strategy都会直接抛出ValueError而on为None时则要求left_on与right_on必须同时给出。1.2 Semi 与 Anti Join最实用的过滤利器Semi 和 Anti Join 本质上就是按另一张表的存在性过滤常用于去重和黑名单过滤。原文档给出的反连接示例——从主表中删除所有id出现在ids_to_remove中的行# Remove all rows from df whose id appears in ids_to_remove df_filtered df.join(ids_to_remove, onid, howanti)在底层实现中semi / anti Join 被专门处理src/daft-local-execution/src/join/anti_semi_join.rs中的probe_anti_semi会按is_semi标志决定保留匹配行semi还是未匹配行anti并且 hash_join.rs 的注释 显示其支持左侧构建 位图bitmap的优化路径来追踪匹配状态。一个常被忽略的性能点semi / anti Join 只输出左侧表的列不需要携带右侧任何数据列因此在大表 vs 小 ID 集过滤场景中它们天然比 left/inner Join 更省内存——这也是下文策略选择中反复推荐它们的原因。二、列名冲突prefix与suffix当左右两侧存在非连接键的同名列时Daft 会为右侧冲突列加上right.前缀。可通过prefix前缀与suffix后缀参数自定义二者可以同时使用# Default: conflicting columns get right. prefix joined df1.join(df2, onkey) # Schema: key, value, right.value # Custom suffix instead joined df1.join(df2, onkey, suffix_other) # Schema: key, value, value_other这与源码 docstringdataframe.py#L3934-L3935及参数默认值完全一致prefix默认right.、suffix默认。实际效果可参考源码中的官方示例dataframe.py#L3962-L3974 df1 daft.from_pydict({a: [w, x, y], b: [1, 2, 3]}) df2 daft.from_pydict({a: [x, y, z], b: [20, 30, 40]}) joined_df df1.join(df2, left_ondf1[a], right_ondf2[a]) joined_df.show() ╭────────┬───────┬─────────╮ │ a ┆ b ┆ right.b │ │ --- ┆ --- ┆ --- │ │ String ┆ Int64 ┆ Int64 │ ╞════════╪═══════╪═════════╡ │ x ┆ 2 ┆ 20 │ ├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌┤ │ y ┆ 3 ┆ 30 │ ╰────────┴───────┴─────────╯ (Showing first 2 of 2 rows)三、三种执行策略Hash / Broadcast / Sort-Merge默认strategyNone时Daft 的查询优化器会自动挑选策略当你知道优化器不知道的信息例如这张表其实非常小时可以显式覆盖。三种策略在 join.rs 的JoinStrategy枚举 中定义pub enum JoinStrategy { Hash, // 通用适用于所有 Join 类型 SortMerge, // 仅支持 Inner Broadcast, // 不支持 Outer // 以下为内部使用用户不可指定 Cross, KeyFiltering, }注意Cross与KeyFiltering是仅供内部使用的策略用户传入字符串时只能使用hash/sort_merge/broadcastjoin.rs#L139-L150 会拒绝其他取值并抛出TypeError。3.1 Hash Join默认策略通用兜底两侧表都按连接键做 Hash 分区hash-partitioned并共置co-located然后进行哈希匹配。这是唯一能覆盖所有 Join 类型的通用策略df.join(other, onkey, strategyhash)代价是两侧都要被 Shuffle因此对于大数据集来说它是内存最吃紧的策略。如果遇到内存压力参见 Managing Memory Usage并评估改用 broadcast 或 sort_merge。3.2 Broadcast Join小表广播零 Shuffle将一侧表复制到每个 Worker 上另一侧完全不需要 Shuffle——当一张表小到每个 Worker 内存都能放下时代价显著低于 Hash Join# Broadcast the small lookup table to all workers df_large.join(df_small, onkey, strategybroadcast)哪一侧被广播取决于 Join 类型Join TypeBroadcast SideinnerThe smaller table (auto-selected)leftRight tablerightLeft tablesemi,antiThe filtering side (right table)限制Broadcast Join 不支持 outer全外连接。这一限制在 Python 层就有强制校验dataframe.py#L3999-L4000elif join_strategy JoinStrategy.Broadcast and join_type JoinType.Outer: raise ValueError(Broadcast join does not support outer joins)3.3 Sort-Merge Join按序归并单趟线性扫描两侧先按连接键排序再进行一趟线性归并。适合数据本身已按连接键有序的场景或在内存压力下替代 Hash Join 的昂贵 Shuffledf.join(other, onkey, strategysort_merge)限制Sort-Merge 只支持 Inner Join。同样在 Python 层有校验dataframe.py#L3997-L3998if join_strategy JoinStrategy.SortMerge and join_type ! JoinType.Inner: raise ValueError(Sort merge join only supports inner joins)3.4 策略 × 类型支持矩阵速查策略支持的 Join 类型是否 Shuffle适用场景hash默认全部七种两侧都 Shuffle通用兜底任何规模broadcast除 outer 外全部仅一侧小表广播另一侧零 Shuffle一侧小表维表、ID 过滤集、查找表sort_merge仅 inner排序而非哈希分区数据已预排序、内存受限的等值连接四、如何选择策略从strategyNone出发的决策路径对绝大多数工作负载让strategyNone由优化器决策是正确的选择。显式覆盖只在以下三种场景一侧远小于另一侧查找表、过滤集、维表使用strategybroadcast避免全量 Shuffle。这在去重流水线中尤其常见——用大表去 Join 一个小的 ID 集合来剔除记录两侧都很大且逼近内存上限考虑是否能改造成 semi / anti Join它们会提前丢弃不需要的列通过df.repartition()增加分区数或设置DAFT_SHUFFLE_ALGORITHMflight_shuffle启用磁盘溢写数据已按连接键预排序strategysort_merge对 Inner Join 可以完全跳过分区步骤。关于第 2 点中的repartition()有一个容易踩坑的细节源码 dataframe.py#L3835-L3848 明确提示NativeRunner本地原生执行器下repartition是不支持的空操作no-op并会给出警告需要改用 RayRunnerdaft.set_runner_ray()才能真正生效。此外repartition(n)不带列参数时执行的是随机重分区全局 Shuffle属于昂贵的操作若只是想调整分区数量而不想重排数据应优先考虑更廉价的df.into_partitions(n)。五、分布式 Join 与内存调优Ray 场景在 Ray 上运行时所有涉及 Shuffle 的 Join 都受制于 Ray Object Store 的内存上限。如果连接键数据放不进分布式内存按以下顺序尝试1. 先试 Broadcast只要一侧足够小就彻底绕开 Shuffle。2. 启用 Flight Shuffle 溢写对大表 × 大表 Join设置DAFT_SHUFFLE_ALGORITHMflight_shuffle并指向本地磁盘卷用于溢写daft.set_execution_config(flight_shuffle_dirs[/mnt/spill])这些配置项在源码 daft-config 的DaftExecutionConfig中有明确定义与默认值shuffle_algorithm默认auto由系统自动选择算法可选pre_shuffle_merge、map_reduce等见 config 测试flight_shuffle_dirsFlight Shuffle 的溢写目录列表默认[/tmp]flight_shuffle_compression溢写数据压缩算法默认lz4。也就是说如果你不显式配置flight_shuffle_dirs溢写会落到/tmp——在分布式环境中这可能不是理想位置务必显式指向持久化本地卷。3. 增加分区数在 Join 之前插入df.repartition(n)降低单个分区的内存压力。注意上一节提到的 NativeRunner 限制。更多内存管理细节参见 Managing Memory Usage 与 Partitioning and Batching。六、Join 排序多表连接时优化器如何重排当一个查询涉及三张及以上的表时优化器会重排 Join 图以最小化中间结果的大小。该逻辑实现在优化规则 ReorderJoins 中默认使用暴力枚举brute-force其上限为7 个关系/// Maximum number of relations for brute force (O(n!) enumeration). const BRUTE_FORCE_MAX_RELATIONS: usize 7; /// Maximum number of relations for DP-ccp (O(3^n) enumeration). const DP_CCP_MAX_RELATIONS: usize 12;对于更大的 Join 图可启用实验性的DP-ccp 枚举器基于 Moerkotte Neumann 2006 的论文设置环境变量export DAFT_DEV_ENABLE_DP_CCP_JOIN_ORDERING1该环境变量在 daft-config 中被解析并写入DaftPlanningConfig.enable_dp_ccp_join_ordering随后在优化器中切换枚举器实现mod.rs#L55-L59let join_order: join_graph::JoinOrderTree if self.use_dp_ccp { dp_ccp_join_order::DpCcpJoinOrderer {}.order(join_graph) } else { brute_force_join_order::BruteForceJoinOrderer {}.order(join_graph) };启用 DP-ccp 后关系数上限从 7 提升到12枚举复杂度从 O(n!) 降到 O(3^n)。两点需要谨慎对待在不超过 7 个关系的 Join 图上DP-ccp 与暴力枚举产生完全相同的计划在更大的图上扩张的搜索空间可能暴露出当前成本模型的弱点从而产生次优计划。官方文档说明该限制将随着统计信息与成本估计的完善而解决在 Daft 仓库的 issue 跟踪中推进。七、结语Daft 的 Join 体系可以概括为七种类型 × 三种策略类型决定输出什么策略决定怎么算。选型时记住三条主线——小表配 broadcast、预排序数据配 sort_merge、其余交给优化器默认的 hash在大数据量 Ray 分布式场景下用 semi/anti 提前瘦身、用 flight shuffle 溢写兜底、用 repartition 分散压力。若想深入执行细节可以从 join.rs 的类型系统、dataframe.py 的参数校验与 docstring、以及 reorder_joins 的枚举器实现继续探索。【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考