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]

Reply via email to