kaxil commented on PR #72938:
URL: https://github.com/apache/airflow/pull/72938#issuecomment-5761917718

   End-to-end runs against the real OpenAI Batch API and the real Anthropic 
Message Batches API, on Airflow `main` in breeze with the deferrable path 
(worker submit, triggerer poll, worker landing). Same Dag shape for both 
providers: three batch tasks plus a downstream task that reads the JSONL rows 
back through `ObjectStoragePath`.
   
   ## OpenAI
   
   | Task | Shape | Proves |
   |---|---|---|
   | `classify_reviews` | `LLMBatchOperator`, `gpt-4.1-mini`, 
`output_type=Sentiment` (Pydantic), `request_params={"temperature": 0, "user": 
...}` | Structured output via `response_format`, OpenAI body params pass 
through |
   | `summarize_reviews` | `@task.llm_batch`, model taken from the connection's 
Model field, per-request `system_prompt` / `max_tokens` / `params` overrides | 
Decorator path, per-request overrides (one request answered in French as 
instructed) |
   | `extract_keywords` | `gpt-5-mini`, `output_type=list[str]`, 
`request_params={"reasoning_effort": "low"}`, `fail_on_partial_error=True` | 
Reasoning-model parameter, `max_completion_tokens` translation, non-object 
schema wrapping/unwrapping |
   
   All 12 requests came back `success`; every manifest reconciled 
(`request_count == sum(counts)`, `terminal_reason: succeeded`).
   
   ![Dag run with all tasks 
green](https://github.com/user-attachments/assets/134d00cc-d12b-400b-a668-67bb61922fc4)
   
   Task log for the operator: submit on the worker, one deferral, trigger 
fires, results landed with per-bucket counts.
   
   ![classify_reviews 
log](https://github.com/user-attachments/assets/b2ff6269-7b4b-48fb-998e-b77d37ee6b33)
   
   XCom for the `gpt-5-mini` task: the `batch_id` key pushed at submit time and 
the manifest as the return value.
   
   ![extract_keywords XCom 
manifest](https://github.com/user-attachments/assets/80bb2009-4ebb-4237-8427-28111a30fc1e)
   
   Clearing a finished task re-attaches to the recorded batch and re-lands the 
same file in a few seconds, with no new submission:
   
   ![re-attach on 
clear](https://github.com/user-attachments/assets/2c086ffd-265a-4cf6-a25d-2d4358b206b6)
   
   Sample rows from the landed JSONL (`gpt-5-mini`, `list[str]` output):
   
   ```json
   {"custom_id": "4fdf6648f4efda69-0", "index": 0, "status": "success", 
"output": ["loved", "would buy again", "positive"], "raw_output": null, 
"error": null, "model": "gpt-5-mini-2025-08-07", "usage": {"input_tokens": 61, 
"output_tokens": 89}, "finish_reason": "stop"}
   {"custom_id": "4fdf6648f4efda69-1", "index": 1, "status": "success", 
"output": ["broke", "one-use", "disappointed"], "raw_output": null, "error": 
null, "model": "gpt-5-mini-2025-08-07", "usage": {"input_tokens": 62, 
"output_tokens": 89}, "finish_reason": "stop"}
   ```
   
   ## Anthropic
   
   | Task | Shape | Proves |
   |---|---|---|
   | `classify_reviews` | `LLMBatchOperator`, `claude-haiku-4-5`, 
`output_type=Sentiment`, `request_params={"temperature": 0.2, "metadata": 
{"user_id": ...}}` | Structured output via a forced tool (`finish_reason: 
tool_use`), Anthropic Messages body params pass through |
   | `summarize_reviews` | `@task.llm_batch`, model from the connection, one 
request overriding `model` to `claude-sonnet-4-5` plus `system_prompt` / 
`max_tokens` / `params` overrides | Decorator path, per-request model override 
inside one batch (Anthropic allows it, OpenAI does not), French override 
honoured |
   | `extract_keywords` | `claude-haiku-4-5`, `output_type=list[str]`, 
`fail_on_partial_error=True` | Non-object schema wrapping/unwrapping through 
the tool input |
   
   All 12 requests `success`, manifests reconciled. Anthropic results arrive 
out of order; rows carry their `index` and the manifest says `ordered: false`, 
`rejoin_key: index`.
   
   ![Anthropic Dag run with all tasks 
green](https://github.com/user-attachments/assets/ce211591-cbcd-4039-beba-8c73c19ce740)
   
   Triggerer polling and landing for the mixed-model decorator task:
   
   ![summarize_reviews 
log](https://github.com/user-attachments/assets/2bf62738-e5be-4dc1-8c76-ddc073c2a583)
   
   ![classify_reviews XCom 
manifest](https://github.com/user-attachments/assets/39fc15cf-0c8e-493f-ac9b-eaadd7d35ebf)
   
   Rows from the mixed-model batch, showing the per-request model override 
landing on Sonnet while the rest ran on Haiku:
   
   ```json
   {"index": 0, "status": "success", "output": "The reviewer highly enjoyed the 
product and plans to purchase it again in the future.", "model": 
"claude-haiku-4-5-20251001", "usage": {"input_tokens": 30, "output_tokens": 
19}, "finish_reason": "end_turn"}
   {"index": 1, "status": "success", "output": "Cet article s'est cassé après 
une seule utilisation, ce qui est très décevant.", "model": 
"claude-haiku-4-5-20251001", "usage": {"input_tokens": 33, "output_tokens": 
27}, "finish_reason": "end_turn"}
   {"index": 2, "status": "success", "output": "Product works as advertised but 
is unremarkable.", "model": "claude-sonnet-4-5-20250929", "usage": 
{"input_tokens": 32, "output_tokens": 14}, "finish_reason": "end_turn"}
   ```
   
   ## Timing
   
   | Provider | Requests | Submit to landed |
   |---|---|---|
   | OpenAI | 3 to 6 | 1 to 4 minutes |
   | Anthropic | 3 | 1 minute (keywords), 19 minutes (mixed-model summaries) |
   | Anthropic | 6 | 21 minutes |
   
   ## Changes that came out of the runs
   
   Both are in the latest commit. The operator now logs a line when it 
re-attaches to a recorded batch and another when it lands results; before this 
a successful landing was silent in the task log. The operator guide gained a 
"Provider-specific parameters" section with a per-provider translation table 
for `system_prompt`, `max_tokens`, `output_type`, `request_params` and 
`completion_window`, and the example Dags were rewritten into four: basic 
operator with a downstream reader, decorator, provider params for both 
providers side by side, and per-request overrides.
   
   One thing to note from the Anthropic log above: every status poll shows up 
as an INFO line (`HTTP Request: GET .../batches/...`), so a long batch fills 
the triggerer log with one line per `poll_interval`. That line is httpx's own 
request log at INFO; the OpenAI task logs above do not show it. Not changed 
here; worth a look as a follow-up.
   


-- 
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]

Reply via email to