Databases 25 min read

StarRocks AI Function: Embedding LLM Inference into the SQL Execution Pipeline

StarRocks introduces AI Functions that let SQL directly call LLMs for generation, classification, extraction, and embedding, integrating model inference into the query optimizer and pipeline execution engine with async dual-pipeline architecture, backpressure, and cost-based optimization to reduce unnecessary calls and manage latency.

StarRocks
StarRocks
StarRocks
StarRocks AI Function: Embedding LLM Inference into the SQL Execution Pipeline

In modern data analysis, database queries are no longer the endpoint; retrieved text or images often need LLM-based summarization, classification, extraction, semantic judgment, or vectorization before further analysis. The common pattern — query data via SQL, then call models from the application layer — works for small volumes but breaks down when model calls become routine: the application must handle data movement, task splitting, concurrency, rate limiting, retries, and result write-back, while database governance (filtering, permissions, resource management, observability) cannot cover the external inference path.

StarRocks addresses this by introducing AI Functions, allowing users to invoke model capabilities directly in SQL. The functions cover four categories: generation and transformation ( ai_complete, ai_summarize, ai_translate), judgment and structured processing ( ai_filter, ai_classify, ai_extract, ai_similarity), vectorization ( ai_embed, ai_embed_multimodal), and multi-row inference ( ai_agg, ai_agg_summary). Model outputs appear as data columns that can further participate in Filter, Join, aggregation, or table writes. The interface draws inspiration from Snowflake Cortex AI Functions but adapts to StarRocks' vectorized MPP engine.

FE: Defining the Execution Boundary with AIProject

Remote model calls break two assumptions of ordinary scalar UDFs: computation is local, and a chunk yields results in one short call. A chunk may contain thousands of rows; synchronous per-row waits kill throughput, while unbounded concurrency explodes memory. The cost model also changes: predicates depending on AI results can only run after inference, TopN on model scores cannot be pushed early, and duplicate identical AI expressions would cause redundant calls and extra cost.

Therefore, StarRocks does not leave AI Functions as ordinary expressions in Project. The Frontend (FE) extracts model calls into a dedicated AIProject logical operator, giving remote inference a distinct data entry, concurrency window, and lifecycle. During analysis, the FE resolves function overloads, validates model names, classification sets, return formats, and parameters, and checks USAGE permissions on AI Model Resources. The optimizer then bottom-up extracts AI calls from the expression tree, representing each result with a new ColumnRef. Nested calls become multiple LogicalAIProject layers; identical calls at the same level share a ColumnRef to evaluate once per plan (not per distinct text value). AI aggregates ( ai_agg, ai_agg_summary) take a different path: array_agg collects group data, then an internal AI scalar function processes the array. For small arrays, the whole input fits in one model context (Stuff mode); larger inputs are split by AIArrayDispatcher into Map/Reduce rounds, all reusing the unified scheduling, rate-limiting, and retry machinery.

With AIProject as a distinct boundary, the optimizer can reorder operators without changing semantics: deterministic predicates on base columns push before AI calls; TopN without OFFSET whose sort key passes through AIProject can be duplicated above it to trim candidates early (the final TopN remains to guarantee correctness); ordinary LIMIT pushes through row-count-preserving projections but stops before filters that depend on AI results. This turns "which rows need model calls" into an optimizer problem — scan and local compute shrink the candidate set first, AIProject processes only rows that truly need inference, and downstream operators consume model outputs. For high-cost inference, eliminating one unnecessary call outweighs optimizing a few CPU instructions.

In physical planning, LogicalAIProject becomes PhysicalAIProject and then AIProjectNode, carrying output expressions, common expressions, and model-capability-organized configs. The Backend (BE) parses these and mandates Pipeline-mode execution, preventing fallback to synchronous get_next().

BE: Dual Pipeline Decouples Fast Data from Slow Service

AIProjectNode::decompose_to_pipeline()

splits one logical node into two pipelines. The upstream pipeline (Scan, Filter, Join, etc.) ends with AIBufferSinkOperator; the downstream starts with AISourceOperator. They connect via a fragment-shared AIChunkBuffer — a multi-producer, multi-consumer FIFO queue, not per-driver channels. Upstream drivers write full chunks; downstream sources compete for chunks, enabling concurrent remote inference. Order is not guaranteed; explicit ORDER BY is required for ordered results. AIChunkBuffer absorbs short-term rate differences between local production and remote inference, monitoring chunk count and memory watermark. When either threshold is hit, Sink::need_input() returns false, putting the driver into OUTPUT_FULL and yielding the thread. After a source consumes data, PipeObservable wakes the upstream. A single oversized chunk is admitted if the queue is empty, avoiding starvation. Each sink decrements a producer counter; only when all sinks finish does the buffer enter EOS. Conversely, if LIMIT is satisfied or the query cancels, sources close the buffer, discard unconsumed chunks, and wake waiting producers — a bidirectional termination protocol that avoids deadlocks on normal, early, or cancelled completion.

Sources split each chunk into smaller sub-chunks per ai_function_sub_chunk_size. Splitting happens on the consumer side: the buffer still uses StarRocks' native chunk as the handoff unit, while the AI execution side uses smaller tasks to control per-task data volume and concurrency granularity. For plans with LIMIT, sources share an atomic row budget, reserving rows when splitting sub-chunks to prevent parallel instances from collectively overshooting the limit.

Each async sub-chunk task clones an independent expression context and holds a lifecycle reference to QueryContext. Common expressions and regular projections run under the query's memory management; AI expressions evaluate only through AISourceOperator 's dedicated entry, avoiding shared mutable FunctionContext across concurrent tasks and ensuring RuntimeState stays alive until the async task finishes.

Async Execution: Driver Doesn't Wait for Network, bthread Waits for Completion

The dual pipeline defines the data-flow boundary; the async task bridge in AISourceOperator moves the remote wait off the Pipeline driver. For a sub-chunk, AISourceOperator submits a Setup Task to the WorkGroup's ScanExecutor. The Setup Task launches a bthread and returns immediately, freeing the ScanExecutor's pthread. Inside the bthread, the AI expression builds the request and calls AITaskDispatcher. While waiting for HTTP response, QPS tokens, or retry windows, bthread::ConditionVariable parks the bthread and yields the underlying pthread, keeping ordinary Pipeline threads unblocked.

For Chat-type functions (generation, classification, extraction), each non-null row in a sub-chunk creates an AITask carrying the original row index; results are written back by index to produce an aligned result column regardless of completion order. Embedding functions add a micro-batch layer: multiple non-null texts merge into one vector request; returned vectors are written back by original position. Sub-chunk controls engine-side task size; micro-batch controls provider-side request batching — two distinct knobs.

Outbound requests pass two process-level admission controls: (1) a QPS token bucket partitioned by endpoint, credentials, and model capability; (2) a global Inflight limit shared by the BE. If a request gets a QPS token but cannot acquire an Inflight slot, the dispatcher returns the token and retries later, preventing phantom requests from consuming rate budget. After admission, AiHttpClient hands an immutable request to a process-level libcurl multi transport maintained by a dedicated network thread (no expression computation). On HTTP completion, a callback updates state and wakes the bthread; the provider adapter parses the response, extracting token usage and errors.

When a sub-chunk finishes, the bthread forces a Collect Task onto the ScanExecutor to write the result chunk into the source's result queue and wake the waiting driver. Failures follow the same path: transport errors, retryable HTTP codes, and 429 responses retry within a configured limit; 429 triggers exponential backoff on the corresponding rate-limit bucket, and each retry re-passes both admission controls. Request timeout is bounded by the query's remaining execution time; cancellation signals cover rate-limit wait, network transfer, and result collection. Row-level errors can either abort the query or set the row's result to NULL per session policy; query cancellation or deadline expiry are not downgraded to row errors.

Protocol differences across models are confined to the Provider layer. The plan declares the required model capability; the Provider handles request encoding and response parsing; scheduling, rate-limiting, memory management, and metrics collection are shared. Adding a new model protocol does not require re-implementing the Pipeline execution machinery.

Backpressure: An Observable Feedback Chain, Not a Single Switch

AI query pressure can originate from three directions: upstream scanning too fast, model endpoint slowing down, or downstream consumption lagging. StarRocks employs multiple continuous boundaries rather than a single concurrency knob. Buffer watermark limits upstream pre-production; Source limits concurrent I/O tasks and result queue size; process-level QPS and Inflight caps constrain actual outbound HTTP requests.

When model response slows, sub-chunk processing time grows, Source gradually stops pulling from Buffer; Buffer hits its watermark, upstream drivers throttle production. Downstream backpressure propagates similarly through the result queue up to the scan layer. The Query Profile exposes this chain with concrete metrics: PeakAIBufferChunks, PeakAIBufferBytes, AIBufferBackpressureCount at the pipeline boundary; PeakIOTasks and PeakResultQueueSize on the Source side; AIQpsWaitTime and AIInflightWaitTime for admission waits; AIHttpRttTime and AIRetryCount for endpoint and network health. These are cumulative across concurrent tasks, so they can exceed wall-clock time and must not be summed as total latency. Diagnosing slow queries means identifying the dominant wait layer: was the candidate set sufficiently reduced before inference? Is Buffer persistently backpressured? Are QPS or Inflight the bottleneck? Or does model RTT exhibit long tails?

Async execution does not escape memory governance. Chunks and expression results are tracked by the instance-level memory tracker; request bodies and retained response bodies are explicitly charged to the query's lifecycle, with a separate cap on response body size.

Each BE independently maintains its token buckets, Inflight limits, and HTTP transport. If a query spans multiple BEs sharing the same model service quota, the operator must plan rate-limit configs based on the aggregate potential request volume across all nodes, or delegate global quota enforcement to a model gateway.

Conclusion

StarRocks AI Function bridges two systems with fundamentally different runtime characteristics: a high-throughput, vectorized data execution engine on one side, and a high-latency, quota-constrained model service on the other. In the FE, AIProject establishes a clear execution boundary for model calls and pushes filtering and candidate trimming before inference. In the BE, dual pipelines with a bounded buffer decouple local data production from remote inference; sub-chunks, ScanExecutor, and bthreads carry async tasks; AITaskDispatcher, process-level rate limiters, and libcurl transport manage outbound requests; and Query Profile ties data backpressure, wait times, call volumes, and token costs back to the originating SQL. To the user, AI Functions remain composable SQL functions alongside Join, aggregation, and DML; to the engine, model calls become first-class remote workloads with explicit boundaries, backpressure, cancellation, observability, and metering — a fundamental difference from "making an HTTP request inside a database function."

Related code has been merged into the StarRocks Main branch and will ship in the next release. Follow-up articles will cover usage patterns and scenarios for each capability as versions progress.

Original Source

Signed-in readers can open the original source through BestHub's protected redirect.

Sign in to view source
Republication Notice

This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactadmin@besthub.devand we will review it promptly.

StarRocksquery optimizationLLM Inferencebackpressureasync processingpipeline executionAI FunctionSQL execution engine
StarRocks
Written by

StarRocks

StarRocks is an open‑source project under the Linux Foundation, focused on building a high‑performance, scalable analytical database that enables enterprises to create an efficient, unified lake‑house paradigm. It is widely used across many industries worldwide, helping numerous companies enhance their data analytics capabilities.

0 followers
Reader feedback

How this landed with the community

Sign in to like

Rate this article

Was this worth your time?

Sign in to rate
Discussion

0 Comments

Thoughtful readers leave field notes, pushback, and hard-won operational detail here.