ashoka1981 opened a new pull request, #68666:
URL: https://github.com/apache/doris/pull/68666
### What problem does this PR solve?
Issue Number: close #68626
Related PR: #62494 (introduced the per-block batching this builds on)
Problem Summary:
The scalar AI functions (`ai_filter`, `ai_classify`, `ai_sentiment`,
`ai_summarize`, ... — everything built on `AIFunction<Derived>::execute`)
already batch rows into one provider request per `ai_context_window_size`
bytes, but the batches of a block are issued strictly one after another:
`execute_impl` builds a batch, blocks on the HTTP round trip, appends the
results, and only then builds the next batch.
For the common query shape
```sql
SELECT id, review_text, ai_filter('res', CONCAT('...', review_text, '...'))
FROM t;
```
the planner produces a single fragment whose scan projection feeds the
result sink directly, i.e. one instance (`Total Instances Num: 1` in the
profile). Combined with the sequential loop, the whole query has exactly one
provider request in flight at any time, regardless of
`parallel_pipeline_task_num`, tablet count or core count. Query time is the sum
of every batch's provider latency. (The `WHERE ai_filter(...)` shape fares
better because the predicate runs inside the scan node with one instance per
scanner, but per-instance execution is still serial.)
Measured on a single-BE cluster (10 cores), 20,000 rows, each batch one HTTP
request to the provider:
| provider (batch) | before: serial | after: `ai.max_concurrency=4` | notes |
|---|---|---|---|
| GPT-4.1-mini via Azure OpenAI (22 rows/batch) | 2,027.7 s | 491.0 s (4.1x)
| provider latency unchanged (p50 2.1 s), 0 rate-limit errors |
| local batching proxy in front of a judgment service (~750 rows/batch) |
24.9 s | 13.8 s | provider-side limit reached |
| same, 100,000 rows (1 MB window) | 96.1 s | 49.9 s | |
The fix: cut a block into its batches first, then execute them through a
bounded, streaming scheduler that keeps up to `ai.max_concurrency` requests in
flight and submits the next batch the moment any one returns. Results are still
appended in batch order, so output is identical row for row. The default (`1`)
is exactly the previous sequential loop.
Design decisions (and the alternatives considered):
1. **Where the parallelism lives: inside `execute()`, not in the planner.**
The root cause is plan shape (blockable projection in a one-instance fragment).
Teaching Nereids to place a blockable projection in a parallel fragment behind
a gather exchange would fix that shape, but it is a much larger change, touches
every plan with an AI function, and still leaves per-instance execution serial.
A bounded pool inside the function works for every plan shape today, is fully
contained in `ai_functions.h`, and composes with any future planner change.
Reviewer question worth asking: is a planner-side change wanted as a follow-up?
2. **Knob on the AI resource (`ai.max_concurrency`), not a session
variable.** The right value is a property of the provider endpoint (its rate
limit / connection budget), not of a query: the same query hitting OpenAI and a
local model should not share one number. It also survives across sessions and
can be changed with `ALTER RESOURCE` without restarting anything. Validation:
positive integer; default `1`. Resources persisted before this property exist
keep working: `AIResource.toThrift()` only sets the field when present, and the
BE treats an unset thrift field (0) as 1.
3. **Streaming rather than waves.** A first cut submitted batches in waves
of N and waited for the wave; the slowest call gated each wave, which cost ~15%
at N=4 against GPT-4.1-mini (573 s vs 491 s). The final version uses a
mutex/condvar loop: exactly N in flight, refill on each completion. The first
error stops new submissions and is returned after the in-flight requests drain,
so no request outlives the call.
4. **One process-wide pool, sized by BE config, not a pool per query.** Pool
threads only block on HTTP, so the pool size
(`ai_function_thread_pool_thread_num`, default 64, queue 10240) is a cap on
total in-flight AI requests per BE rather than CPU parallelism. A per-query
pool would give no such bound and would be created/destroyed constantly. If the
pool is full, `submit_func` fails and the batch runs inline on the calling
thread, so progress is never blocked on pool capacity. Same construction
pattern as `SendBatchThreadPool` in `exec_env_init.cpp`.
5. **Memory accounting.** Pool threads are orphans; each task attaches the
query's `ResourceContext` (`AttachTask`) so request/response buffers are
charged to the query, following `snii_doris_adapter.cpp`'s concurrent reads.
6. **Not serialized.** `AIResource::serialize()/deserialize()` ship `ai_agg`
state between BEs; `max_concurrency` is irrelevant there (`ai_agg` has its own
request path) and adding a field to that stream would break rolling upgrades,
so it is deliberately excluded.
7. **Failure semantics are unchanged.** Any batch failure fails the query,
as before; the only difference is that up to N-1 sibling requests may be in
flight when it does. Query timeout still bounds each request through
`HttpClient::set_timeout_ms(remaining_query_time)`.
Scope: `AIFunction<Derived>::execute` (all scalar AI functions that use the
shared batch loop). `embed` and `ai_agg` have their own request paths and are
not changed.
### Release note
Add the `ai.max_concurrency` AI resource property (default `1`): the maximum
number of provider requests one execution instance keeps in flight for the
scalar AI functions (`ai_filter`, `ai_classify`, `ai_sentiment`, ...). Values
above 1 overlap the batch requests of a block instead of issuing them one at a
time. Adds BE config `ai_function_thread_pool_thread_num` (default 64) and
`ai_function_thread_pool_queue_size` (default 10240) for the shared pool that
runs them.
### Check List (For Author)
- Test
- [x] Regression test: `ai_p0/test_create_ai_resource` —
`ai.max_concurrency` accepted, non-positive value rejected.
- [x] Unit Test:
`AIFunctionTest.ExecuteBatchesConcurrentKeepsBatchOrder`,
`AIFunctionTest.ExecuteBatchesConcurrentPropagatesBatchError` (BE;
`run-be-ut.sh --run --filter='AIFunctionTest.*'` → 53/53 pass);
`AIResourceTest` (FE, 31/31: default and explicit value reach thrift).
- [x] Manual test: single-BE cluster, `ai_filter` over 20K/100K rows
against Azure OpenAI and a local OpenAI-compatible proxy; results
byte-identical to the sequential run (confusion matrix vs ground truth
unchanged), timings in the table above. `SHOW PROC` / `ALTER RESOURCE` round
trip verified; resource created before the change still loads and runs
sequentially.
- [ ] No need to test or manual test.
- Behavior changed:
- [x] No. Default `ai.max_concurrency=1` keeps the sequential loop;
behavior only changes when a resource opts in.
- Does this need documentation?
- [x] Yes. doris-website PR: <link> — document `ai.max_concurrency`
under the AI resource properties (alongside `ai.max_retries`), and the two BE
config items.
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label
🤖 Generated with [Claude Code](https://claude.com/claude-code)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]