NEAREST BY Join: Scaling Vector Search in Databricks Runtime Databricks built vector search directly into Databricks Runtime as a first-class SQL join called NEAREST BY, replacing its earlier VECTOR_SEARCH implementation that federated requests to an external real-time endpoint and was capped by endpoint sizing rather than cluster size. The new stack lowers the top-k ranking join onto SIMD-accelerated distance functions, a bounded top-k aggregate, and a fused Photon operator with a custom GEMM kernel, plus an optional IVF index stored as an ordinary liquid-clustered Delta table so APPROX queries score a fraction of the base vectors. Databricks cites batch workloads such as a payments company matching 100M+ daily transactions against 140M merchant embeddings and a quantitative fund running million-query batches against a 50M-vector corpus as the target use cases. How we built vector search into Databricks as a first-class SQL join, with deep kernel optimizations in Photon and a vector index in an open storage format. by Zero Qu Zero https://www.databricks.com/blog/author/zero-qu-zero , Alexis Schlomer https://www.databricks.com/blog/author/alexis-schlomer , Akash Nayar https://www.databricks.com/blog/author/akash-nayar , Yingyi Bu https://www.databricks.com/blog/author/yingyi-bu and Sergei Tsarev https://www.databricks.com/blog/author/sergei-tsarev 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. 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. Our first version of the VECTOR SEARCH https://docs.databricks.com/aws/en/sql/language-manual/functions/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. This resulted in a deliberately small yet deep stack: a new join syntax, NEAREST BY https://docs.databricks.com/aws/en/sql/language-manual/sql-ref-syntax-qry-select-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. 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. 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 core building blocks are a family of vector SQL functions over ARRAY