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`).

Task log for the operator: submit on the worker, one deferral, trigger
fires, results landed with per-bucket counts.

XCom for the `gpt-5-mini` task: the `batch_id` key pushed at submit time and
the manifest as the return value.

Clearing a finished task re-attaches to the recorded batch and re-lands the
same file in a few seconds, with no new submission:

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`.

Triggerer polling and landing for the mixed-model decorator task:


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]