OpenAIEmbeddingOperator¶
Use the OpenAIEmbeddingOperator to
interact with the OpenAI API to create embeddings for given text.
Using the Operator¶
The OpenAIEmbeddingOperator requires the input_text as an input to embedding API. Use the conn_id parameter to specify the OpenAI connection to use to
connect to your account.
A single string or token array returns one embedding vector. A list of strings or token arrays returns one vector per input item in the same order.
An example of using the operator:
OpenAIEmbeddingOperator(
task_id="embedding_using_xcom_data",
conn_id="openai_default",
input_text=task_to_store_input_text_in_xcom(),
model="text-embedding-3-small",
)
OpenAIEmbeddingOperator(
task_id="embedding_using_callable",
conn_id="openai_default",
input_text=input_text_callable(
"input_arg1_value",
"input2_value",
input_kwarg1="input_kwarg1_value",
input_kwarg2="input_kwarg2_value",
),
model="text-embedding-3-small",
)
OpenAIEmbeddingOperator(
task_id="embedding_using_text",
conn_id="openai_default",
input_text=texts,
model="text-embedding-3-small",
)
OpenAIResponseOperator¶
Use the OpenAIResponseOperator to generate a
model response with the OpenAI Responses API, OpenAI’s recommended interface for text generation and
tool use. The operator returns the response’s aggregated output text. When do_xcom_push is
enabled (the default), execute also pushes two XCom keys: response_id (the response’s ID,
usable as a downstream task’s previous_response_id for chaining) and usage (the response’s
token usage, or None when the API omits it). usage is the nested dict returned by
ResponseUsage.model_dump(): top-level input_tokens, output_tokens and total_tokens
counts, plus the nested input_tokens_details and output_tokens_details dicts.
input_tokens_details.cached_tokens is part of the input_tokens total, not
additional to it, so pricing a run correctly means reading the breakdown rather than
treating input_tokens as a single uniformly priced count – see OpenAI’s prompt
caching guide for how
cached tokens are priced. Beyond that, usage reports token counts only – OpenAI’s
response carries no cost field, so turning any of these counts into a price means
multiplying by your own per-token rate. When usage is not None it also carries a
try_number key recording which attempt produced it – XCom is cleared at the start of
every attempt, so on a retried task instance the usage XCom only ever reflects the
most recent attempt, and try_number makes that scope explicit instead of letting it
silently under-report total spend across retries. Setting do_xcom_push=False skips both pushes.
It also disables the operator’s own return_value XCom (standard BaseOperator
behavior), so a downstream task reading openai_response.output – which implicitly
reads the return_value key – loses that value too.
Using the Operator¶
The OpenAIResponseOperator requires the input_text prompt. Use the conn_id parameter to
specify the OpenAI connection to use, and response_kwargs to pass through options such as
tools, conversation or previous_response_id. response_kwargs is templated, so
previous_response_id can reference a Dag’s upstream response_id XCom directly. Since
response_kwargs is templated, a literal {{ ... }} value you want sent to the API as-is
(for example inside a prompt’s instructions) must be wrapped in a {% raw %} block, for
example {% raw %}{{ not_a_variable }}{% endraw %}.
Use max_output_tokens and max_tool_calls to cap generation per run – both are templated,
so a ceiling can vary by environment or Dag run without hardcoding it. max_output_tokens caps
the number of tokens generated; max_tool_calls caps the number of built-in tool calls the model
may make. Both limits are enforced by the OpenAI API itself; OpenAI exposes no monetary cost limit
on the Responses API, so this operator has no cost cap. For a monetary limit, use
apache-airflow-providers-common-ai instead. Hitting max_output_tokens does not
fail the request: the response comes back with status="incomplete", so return_value will
not raise – but it is not guaranteed to be truncated text either. A reasoning model can spend
the entire ceiling on reasoning tokens and return an empty output_text, in which case
return_value is an empty string. Hitting max_tool_calls is different: the OpenAI API
silently drops any tool calls beyond the ceiling without changing status or setting
incomplete_details – there is no log warning and no signal in return_value, so a run
truncated by max_tool_calls looks identical to a clean run.
A rendered max_output_tokens or max_tool_calls that is blank or whitespace-only – for
example max_output_tokens="{{ params.tokens | default('', true) }}" when params.tokens is
unset – is treated as “no ceiling for this run” rather than raising. This only applies when the same run does
not also set the corresponding key in response_kwargs: the mutually-exclusive-with-response_kwargs
check happens when the operator is constructed and fires regardless of what the template later
renders to.
openai_response = OpenAIResponseOperator(
task_id="openai_response",
conn_id="openai_default",
input_text="Write a haiku about data pipelines.",
response_kwargs={"instructions": "You are a helpful assistant."},
)
# Chains onto the previous response via its response_id XCom, continuing
# the same conversation without resending prior turns. Referencing the
# upstream XComArg here also creates the task dependency automatically --
# no explicit `>>` is needed.
OpenAIResponseOperator(
task_id="openai_response_follow_up",
conn_id="openai_default",
input_text="Now rewrite it as a limerick.",
response_kwargs={"previous_response_id": XComArg(openai_response, key="response_id")},
)
# ``max_output_tokens`` is templated, so a ceiling can vary by Dag run without hardcoding it.
# This Dag does not declare a ``tokens`` param, so ``params.tokens`` is undefined at render
# time. ``| default('', true)`` renders undefined *and* None values to ``''`` (treated as "no
# ceiling"); the more common ``{{ params.tokens or '' }}`` idiom would instead raise
# ``UndefinedError`` under Airflow's default ``StrictUndefined`` template behavior.
OpenAIResponseOperator(
task_id="openai_response_with_token_ceiling",
conn_id="openai_default",
input_text="Write a haiku about data pipelines.",
max_output_tokens="{{ params.tokens | default('', true) }}",
)
Passing Responses API options¶
See the Responses API reference for the authoritative list
of parameters. response_kwargs passes straight through to the underlying create_response
call, so most keyword arguments the Responses API accepts can be set there, with the exceptions
noted below. What actually works also depends on the openai package version installed in the
environment, not the reference page above: Responses.create accepts no arbitrary keyword
arguments, so passing one the installed package doesn’t recognize raises TypeError before any
request is sent. Use extra_body as a fallback to pass a parameter the installed package doesn’t
know about yet. Options worth knowing about:
background: run the response asynchronously on OpenAI’s side. See the note onbackgroundbelow before using this withOpenAIResponseOperator.stream: return a stream of response events instead of a single completed response. Do not set this onOpenAIResponseOperator:executereadsresponse.statusandresponse.output_text, neither of which exists on the streamed response object, so the task raisesAttributeError. Stream responses from a@taskusingOpenAIHookinstead.store: whether the response is retained on OpenAI’s side, for example so it can later be used as aprevious_response_id. Whendo_xcom_pushis enabled,executepushesresponse.idto theresponse_idXCom regardless ofstore, so a downstream task can retrieve it. Ifstore=False, the pushed id has no practical use: nothing was retained on OpenAI’s side, soprevious_response_idcannot reference it.previous_response_id: the id of a prior response to continue a multi-turn conversation from. Cannot be used together withconversation— pass one or the other, not both.reasoning: configuration for reasoning models, for example{"effort": ...}. The example Dag above (and the operator’s own default) usesgpt-4o-mini, which is not a reasoning model, so this option only takes effect ifmodelis also set to a reasoning model.service_tier: the processing tier the request is served from.prompt_cache_key: an identifier used to route requests to the same prompt cache. How long a cache entry is retained is a separate option whose name depends on the installed package:prompt_cache_retentionat the 2.37.0 floor, deprecated in later releases in favor ofprompt_cache_options.ttl.safety_identifier: a stable identifier for the end user, used for safety and abuse detection.truncation: one of'auto'or'disabled'(the default). Under'disabled', a request whose input exceeds the model’s context window fails with a 400 error;'auto'shortens the input to fit instead.include: additional output fields to include in the response, such as encrypted reasoning content. These fields land onresponse.output, butresponse.output_textonly aggregatesmessage/output_textcontent, so anythingincludeadds is fetched and then discarded byexecute. UseOpenAIHookdirectly to access it.metadata: a mapping of key-value pairs attached to the response for your own bookkeeping.max_output_tokens: an upper bound on the number of tokens the model can generate, including reasoning tokens as well as visible output tokens. Setting this key here (instead of the operator’s ownmax_output_tokensparameter – see above) applies the same validation and type-coercion rules; setting the same ceiling in both places raises when the operator is constructed. The one behavioral difference: a blank value here is popped from the payload sent tocreate_response(rather than never being added, as with the operator argument), and a present key whose value is a literalNonealways raises, since presence of the key – not the value – is what “supplied” means for this path.max_tool_calls: an upper bound on the number of built-in tool calls the model can make. Same validation, coercion, and blank/Nonehandling asmax_output_tokensabove.
Note
OpenAI does not expose a spend or cost ceiling parameter on the Responses API.
max_output_tokens and max_tool_calls are token and call-count limits, not a way to cap
the dollar cost of a run; controlling spend means bounding those counts yourself.
Note
background=True starts the response running asynchronously on OpenAI’s side and returns
before the response finishes. OpenAIResponseOperator is synchronous: it makes one
create_response call and returns response.output_text immediately, so a response
started with background=True comes back incomplete, and the operator logs its own warning
because response.status is not yet "completed". Do not set background=True on
OpenAIResponseOperator. If you need a background response, create it from a @task
using OpenAIHook’s create_response directly.
Using the OpenAIHook for Responses and Conversations¶
The OpenAIHook exposes the Responses and
Conversations APIs directly for use inside @task functions or custom operators:
Responses:
create_response,get_response,delete_responseandcancel_response(the last cancels a response created withbackground=True).Conversations:
create_conversation,get_conversation,update_conversationanddelete_conversation. Pass the conversation id tocreate_response(via the operator’sresponse_kwargsor the hook) to persist state across responses.
For example, to create a conversation and continue it across responses:
hook = OpenAIHook()
conversation = hook.create_conversation()
hook.create_response(input="Hello", conversation=conversation.id)
Note
The Assistants/Threads hook methods (create_assistant, create_thread, create_run and
related) are deprecated, mirroring OpenAI’s deprecation of the Assistants API. Migrate to the
Responses and Conversations methods above.
OpenAITriggerBatchOperator¶
Use the OpenAITriggerBatchOperator to
interact with the OpenAI API to trigger a batch job. This operator is used to trigger a batch job and wait for the job to complete.
Using the Operator¶
The OpenAITriggerBatchOperator requires the prepared batch file as an input to trigger the
batch job. Provide the file_id and the endpoint to trigger the batch job, and use the
conn_id parameter to specify the OpenAI connection to use.
An example of using the operator:
from airflow.providers.openai.operators.openai import OpenAITriggerBatchOperator
batch_id = OpenAITriggerBatchOperator(
task_id="batch_operator_deferred",
conn_id=OPENAI_CONN_ID,
file_id=batch_file_id,
endpoint="/v1/chat/completions",
deferrable=True,
)