跳到主要内容
AI News HubLIVE
来源内容 · 翻译待补全6 分钟阅读

待翻译:NEAREST BY Join: Scaling Vector Search in Databricks Runtime

文章摘要

AI 服务暂时不可用,以下为来源摘要,待恢复后补全翻译:Vector search originated as a serving problem. The classical use case is a chatbot...

待翻译:NEAREST BY Join: Scaling Vector Search in Databricks Runtime
报告错误

纠错通道尚未开通,可先复制下方文章信息留存。

查看更正说明
直接读正文

AI 服务暂时不可用,以下为来源正文,待恢复后补全翻译。

NEAREST BY Join: Scaling Vector Search in Databricks Runtime | Databricks Blog Skip to main content NEAREST BY is a new SQL join for batch vector search: for every query row, find the k nearest rows by vector similarity or distance — exact or approximate. A fused Photon operator with a custom blocked GEMM kernel pushes distance scoring toward the peak arithmetic throughput the hardware offers. Your Lakehouse is also your vector store, with no separate system to sync or operate: the IVF vector index is an ordinary liquid-clustered Delta table that prunes away most partitions on read. Vector search originated as a serving problem. The classical use case is a chatbot or a search bar: one query embedding arrives, and the system is optimized to return the top-k nearest documents within tens of milliseconds. However, a good share of the vector search workloads on our platform are inherently batch-oriented — precomputing exact or approximate nearest neighbors offline rather than looking them up at request time. A payments company matches 100M+ daily transactions against 140M merchant embeddings for entity resolution; a data firm enriches tens of millions of historical records nightly; a quantitative fund runs million-query batches against a 50M-vector corpus for taxonomy tagging. Entity resolution, deduplication, semantic tagging, classification, record enrichment, batch recommendations — these are fundamentally batch workloads: millions of queries against millions to billions of vectors on a schedule, measured by whether the job completes within its SLA at a reasonable cost rather than by the latency of a single lookup. These workloads deserve a very different architecture for better performance, reliability, and cost efficiency — so we went back to first principles. Requirements Aggregate throughput over per-request latency. The unit of success is the entire batch job completing within its SLA at reasonable cost, so the design should trade per-request latency for throughput at every opportunity. Scale on both sides of the join. Up to hundreds of millions of query vectors against billions of base vectors. The system must handle all shapes of query and base cardinality. Elastic parallelism. Batch throughput comes from horizontal scale. The work must partition cleanly across hundreds to thousands of cores, and compute should size itself to the job: scale out for the run, scale down to zero after it. Peak arithmetic per core. Distance scoring is computationally expensive. Scale-out only multiplies what a single core achieves, so the inner loops must run near the theoretical maximum set by the underlying hardware's arithmetic bandwidth (FLOPs/s) and memory bandwidth (bytes/s). Fault tolerance. A job that runs for hours must survive worker loss, transient task failures, and memory pressure through spilling to disk. These are properties of an execution engine, not features we can wrap around a real-time serving endpoint. The Databricks Runtime fits these requirements perfectly — a distributed, fault-tolerant, elastic execution engine built on top of Spark and Photon, a vectorized native C++ query engine. This is exactly why we decided to build vector search directly as an engine native feature rather than relying on separate infrastructure. Architecture Our first version of the VECTOR_SEARCH SQL function was designed to federate requests to an external real-time Vector Search endpoint. It was implemented as a Generate node streaming one query row at a time: every row incurred a network request, a response to deserialize, and possibly retries. It worked, but exposed a performance ceiling — throughput was capped by real-time endpoint sizing rather than runtime cluster size, with the runtime engine reduced to a dispatcher. It also missed the true shape of the query. A batch vector search is not a million small searches. It is one large query: for each row on the left, find the k nearest rows on the right — a top-k ranking join. Executing enormous joins is exactly what the runtime engine excels at. Implementing vector search natively in the runtime engine pays off from two angles. Single copy of the data. embeddings stay in Delta tables on the Lakehouse — no separate vector store, no sync pipeline to keep consistent, no second system to operate and pay for. Single engine for execution. search runs in one engine that scales elastically with the workload, with kernels purposely built for batch query shapes — pushing each core toward peak FLOPs, and let scale-out multiply the rest. The engine already owns fault tolerance: tasks retry automatically, and memory pressure spills to disk. No client-side concurrency control, rate limiting, or retry loops. This resulted in a deliberately small yet deep stack: a new join syntax, NEAREST BY, that makes the top-k ranking join a first-class relational operation; a rewrite that lowers it onto three primitives — SIMD-accelerated distance functions and a bounded top-k aggregate; a fused Photon operator that replaces the plan's entire middle with a custom GEMM kernel; and an optional IVF index built as an ordinary liquid-clustered Delta table, which lets APPROX queries score a fraction of the base vectors with the same kernels. The syntax: a top-k ranking join Existing engines converged on two interface shapes. Postgres with pgvector and Snowflake compose distance operators with ORDER BY … LIMIT — batch then needs a LATERAL subquery per driving row, and the optimizer lacks a pattern to recognize for differentiating KNN and ANN queries. That recognition is also fragile: drift from the expected query shape and the fast path silently disappears. BigQuery exposes a table-valued function — batch is first-class, but column references are strings the parser can't validate. Structurally, batch vector search is a binary relational operation: two table inputs, an output combining both, a per-left-row top-k connecting them. The syntax encodes that structure as a native top-k ranking join: The join is asymmetric, similar to LATERAL: the left side drives, the right side is searched. Ranking direction is explicit: BY SIMILARITY descending, BY DISTANCE ascending. LEFT OUTER keeps query rows with no candidates, and the BY expression is pluggable: any orderable scalar over both sides works, so other scoring expressions can reuse the same clause later. APPROX and EXACT encode a semantic contract. EXACT guarantees the true top-k by exhaustive evaluation; APPROX lets the optimizer substitute an approximate strategy, such as an ANN index, where one applies. So creating or dropping an index can never silently change query results: only queries that specify APPROX consent to approximation. The query rewrite NEAREST BY parses into a logical join node, which the optimizer lowers onto standard relational operators: the rewrite tags each query row with a generated id, scores every (query, base) pair, keeps the k best per id with a grouped top-k, and inlines the kept rows back out: Semantically, this encapsulates the entire feature: a cross join, a scalar scoring expression, and a grouped top-k aggregate. Because every operator is an ordinary relational one, the plan distributes, spills, and retries like any other — correctness and fault tolerance come for free. What the rewrite really isolates is the two primitives all execution time flows through: the distance function that scores a pair, and the aggregate that keeps each group's k best. We implemented every operator in this plan natively in Photon, plus one additional fused operator, purposely built for vector search, that collapses the plan's middle section entirely with a more performant, batch-friendly kernel. The photon kernels The vector functions The core building blocks are a family of vector SQL functions over ARRAY columns. Three of them are responsible for similarity and distance computation: SQL FunctionComputesNearer Means vector_inner_product(a, b)Higher (BY SIMILARITY) vector_cosine_similarity(a, b)Higher (BY SIMILARITY) vector_l2_distance(a, b)Lower (BY DISTANCE) Alongside the similarity and distance functions, we shipped two norm helpers — vector_norm and vector_normalize, and two aggregates, vector_sum and vector_avg. Together they cover both query and index construction: the distance functions score queries and assign rows to their nearest centroid, while the aggregates and normalizers recompute those centroids during k-means. Photon runs the whole family as native SIMD kernels. Every metric is multiply-adds at the core, and a single fused multiply-add (FMA) instruction computes 𝑎 · 𝑏 + 𝑐 across an entire vector register per issue. The kernels are implemented based on four deliberate design choices. Portable by construction. The same kernels must run on every cloud (AWS, Azure, GCP) and every CPU architecture (x86, ARM) we support. Since SIMD width and instruction set differ across them, each kernel is compiled into several ISA-specific clones, and the best one CPU supports is selected at run time (AVX-512 on a modern Intel core, SVE2 on Graviton). Zero-copy input. An ARRAY is contiguous within a row, so each vector is read as a raw pointer into the column's backing buffer with no copy and no per-element indexing inside the hot loop. Relaxed float semantics. Dropping strict IEEE ordering lets the compiler fuse multiply-adds into FMA instructions and split the reduction into parallel accumulator chains, so the loop isn't serialized on one running sum. Nothing scalar in the loop. The hot loops are pure vectorized multiply-add reductions, scalar work like the sqrt and the divide are done once, outside the loop. The top-k aggregate In plain SQL, grouped top-k is a window function: ROW_NUMBER() OVER (PARTITION BY query ORDER BY score). This sorts every partition in full, then throws away everything but the top k ranks. We instead extended the existing max_by / min_by aggregates with a third K parameter overload. The implementation is built around four key properties. One-comparison selection. Aggregation state per group is a bounded heap whose root is the eviction candidate and the admission threshold. Once full, each candidate is accepted or rejected with a single comparison. No sort ever runs over the candidate stream. O(k) state. Memory per group is independent of input size. A query scored against a billion rows only carries k rows of state, and k goes up to 100,000. Late materialization. The heap stores indices, not fully materialized copies. A candidate stays a pointer into the live columnar batch and is copied into aggregate state only if it still survives when the batch is recycled. Distributable. Aggregates have a partial/merge contract. Each partition emits its local top-k, the shuffle moves those k-element arrays instead of raw pairs, and the merge re-inserts them. The result is exact because the global top-k is always a subset of the union of the partials. The roofline model With all the native kernels above, the query plan is fully Photonized, yet still far from optimal at batch scale. The reason is theoretical, not implementational, and the roofline model is a simple and effective way to visualize it. 𝑃𝑎𝑡𝑡𝑎𝑖𝑛𝑎𝑏𝑙𝑒 = 𝑚𝑖𝑛(𝑃𝑝𝑒𝑎𝑘, 𝐴𝐼 × 𝐵𝑊) 𝑃𝑝𝑒𝑎𝑘 is the hardware's peak compute throughput (FLOPs/s), BW is memory bandwidth (bytes/s), and 𝐴𝐼 is the kernel's arithmetic intensity (FLOPs performed per byte moved). Plot 𝑃𝑎𝑡𝑡𝑎𝑖𝑛𝑎𝑏𝑙𝑒 against 𝐴𝐼 and we get the roofline: a diagonal memory roof meeting a horizontal compute roof at the ridge point — the minimum 𝐴𝐼 at which a kernel can be compute-bound. Left to the ridge, only moving fewer bytes per FLOP helps; right of it, the kernel is compute-bound and the arithmetic units themselves are the limit. Concretely, on a reference m6i.2xlarge machine (illustrative, the constants shift with the hardware): Per core Compute throughput roof D [truncated for AI cost control]

展开要点与分析

文章情报

投资人进阶

要点

  • AI 服务暂时不可用,系统已先保留来源内容与降级元数据。
  • Vector search originated as a serving problem. The classical use case is a chatbot...

要点与分析由自动化流程生成,可能有误,请结合原始来源核实。