Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions airflow-core/docs/administration-and-deployment/pools.rst
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,12 @@ descendants.
Note that if tasks are not given a pool, they are assigned to a default pool ``default_pool``, which is
initialized with 128 slots and can be modified through the UI or CLI (but cannot be removed).

Whether deferred tasks occupy pool slots is normally decided per pool via its ``include_deferred`` flag.
A Deployment Manager can instead fix this behavior for the whole cluster with
:ref:`config:core__pool_include_deferred`. When that option is set to ``True`` or ``False``, the configured
value applies to every pool (including pre-existing pools, regardless of their stored flag), and attempts
to explicitly set a conflicting ``include_deferred`` value when creating or updating a pool are rejected.

Using multiple pool slots
-------------------------

Expand Down
7 changes: 5 additions & 2 deletions airflow-core/src/airflow/api/client/local_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -78,10 +78,13 @@ def get_pool(self, name):
pool = Pool.get_pool(pool_name=name)
if not pool:
raise PoolNotFound(f"Pool {name} not found")
return pool.pool, pool.slots, pool.description, pool.include_deferred, pool.team_name
return pool.pool, pool.slots, pool.description, pool.effective_include_deferred, pool.team_name

def get_pools(self):
return [(p.pool, p.slots, p.description, p.include_deferred, p.team_name) for p in Pool.get_pools()]
return [
(p.pool, p.slots, p.description, p.effective_include_deferred, p.team_name)
for p in Pool.get_pools()
]

def create_pool(self, name, slots, description, include_deferred, team_name=None):
if not (name and name.strip()):
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

from airflow.api_fastapi.core_api.base import BaseModel, StrictBaseModel
from airflow.configuration import conf
from airflow.models.pool import Pool


def _call_function(function: Callable[[], int]) -> int:
Expand All @@ -35,6 +36,19 @@ def _call_function(function: Callable[[], int]) -> int:
return function()


def _apply_include_deferred_override(value: bool) -> bool:
override = Pool.get_include_deferred_override()
return value if override is None else override


def _reject_conflicting_include_deferred(value: bool, override: bool) -> None:
if value != override:
raise ValueError(
f"include_deferred is fixed to {override} for all pools by the [core] pool_include_deferred "
"configuration and cannot be set per pool. Please contact your administrator."
)


PoolSlots = Annotated[
int,
Field(ge=-1, description="Number of slots. Use -1 for unlimited."),
Expand All @@ -59,6 +73,9 @@ def _sanitize_open_slots(value) -> int:
class PoolResponse(BasePool):
"""Pool serializer for responses."""

# Report the effective value: the cluster-wide config value takes precedence over the stored column
include_deferred: Annotated[bool, BeforeValidator(_apply_include_deferred_override)]

occupied_slots: Annotated[int, BeforeValidator(_call_function)]
running_slots: Annotated[int, BeforeValidator(_call_function)]
queued_slots: Annotated[int, BeforeValidator(_call_function)]
Expand Down Expand Up @@ -92,6 +109,13 @@ def validate_team_name(self) -> PoolPatchBody:
)
return self

@model_validator(mode="after")
def enforce_include_deferred_override(self) -> PoolPatchBody:
override = Pool.get_include_deferred_override()
if override is not None and self.include_deferred is not None:
_reject_conflicting_include_deferred(self.include_deferred, override)
return self


class PoolBody(BasePool, StrictBaseModel):
"""Pool serializer for post bodies."""
Expand All @@ -108,3 +132,13 @@ def validate_team_name(self) -> PoolBody:
"team_name cannot be set when multi_team mode is disabled. Please contact your administrator."
)
return self

@model_validator(mode="after")
def enforce_include_deferred_override(self) -> PoolBody:
override = Pool.get_include_deferred_override()
if override is None:
return self
if "include_deferred" in self.model_fields_set:
_reject_conflicting_include_deferred(self.include_deferred, override)
self.include_deferred = override
return self
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ class ConfigResponse(BaseModel):
theme: Theme | None
multi_team: bool
rerun_with_latest_version: bool | None = None
pool_include_deferred: bool | None = None

@field_serializer("theme")
def serialize_theme(self, theme: Theme | None) -> dict | None:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2665,6 +2665,11 @@ components:
- type: boolean
- type: 'null'
title: Rerun With Latest Version
pool_include_deferred:
anyOf:
- type: boolean
- type: 'null'
title: Pool Include Deferred
type: object
required:
- fallback_page_limit
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
from airflow.api_fastapi.core_api.security import requires_authenticated
from airflow.configuration import conf
from airflow.models.pool import Pool
from airflow.settings import DASHBOARD_UIALERTS
from airflow.utils.log.log_reader import TaskLogReader

Expand Down Expand Up @@ -69,6 +70,8 @@ def get_configs() -> ConfigResponse:
if conf.has_option("core", "rerun_with_latest_version")
else None
),
# None means the flag is chosen per pool; a boolean means it is fixed cluster-wide.
"pool_include_deferred": Pool.get_include_deferred_override(),
}

config.update({key: value for key, value in additional_config.items()})
Expand Down
21 changes: 21 additions & 0 deletions airflow-core/src/airflow/cli/commands/pool_command.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,11 +27,23 @@
from airflow.cli.simple_table import AirflowConsole
from airflow.cli.utils import deprecated_for_airflowctl
from airflow.exceptions import PoolNotFound
from airflow.models.pool import Pool
from airflow.utils import cli as cli_utils
from airflow.utils.cli import suppress_logs_and_warning
from airflow.utils.providers_configuration_loader import providers_configuration_loaded


def check_include_deferred_choice_allowed(include_deferred: bool) -> str | None:
"""Return an error message when an explicit ``include_deferred`` choice conflicts with the cluster config."""
override = Pool.get_include_deferred_override()
if override is not None and include_deferred != override:
return (
f"include_deferred is fixed to {override} for all pools by the [core] pool_include_deferred "
"configuration and cannot be set per pool."
)
return None


def _show_pools(pools, output):
AirflowConsole().print_as(
data=pools,
Expand Down Expand Up @@ -75,6 +87,9 @@ def pool_get(args):
@providers_configuration_loaded
def pool_set(args):
"""Create new pool with a given name and slots."""
# --include-deferred is a store-true flag, so only a passed flag is an explicit choice
if args.include_deferred and (error := check_include_deferred_choice_allowed(True)):
raise SystemExit(error)
api_client = get_current_api_client()
api_client.create_pool(
name=args.pool,
Expand Down Expand Up @@ -136,6 +151,12 @@ def pool_import_helper(filepath):
failed = []
for k, v in pools_json.items():
if isinstance(v, dict) and "slots" in v and "description" in v:
if "include_deferred" in v and (
error := check_include_deferred_choice_allowed(bool(v["include_deferred"]))
):
print(f"Pool {k}: {error}")
failed.append(k)
continue
pools.append(
api_client.create_pool(
name=k,
Expand Down
12 changes: 12 additions & 0 deletions airflow-core/src/airflow/config_templates/config.yml
Original file line number Diff line number Diff line change
Expand Up @@ -463,6 +463,18 @@ core:
type: integer
example: ~
default: "128"
pool_include_deferred:
description: |
Cluster-wide setting for the ``include_deferred`` flag of pools. When left empty (the default),
each pool keeps its own ``include_deferred`` value, configurable per pool via the UI, API or CLI.
When set to ``True`` or ``False``, the configured value is used for **every** pool (including
pre-existing pools, whatever their stored value) when calculating occupied slots, and users can
no longer choose the flag per pool: attempts to explicitly set a conflicting ``include_deferred``
value when creating or updating a pool are rejected.
version_added: 3.4.0
type: string
example: "True"
default: ""
max_map_length:
description: |
The maximum list/dict length an XCom can push to trigger task mapping. If the pushed list/dict has a
Expand Down
31 changes: 29 additions & 2 deletions airflow-core/src/airflow/models/pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,26 @@ class Pool(Base):
def __repr__(self):
return str(self.pool)

@staticmethod
def get_include_deferred_override() -> bool | None:
"""
Get the cluster-wide ``include_deferred`` value fixed via config, if any.

When ``[core] pool_include_deferred`` is set, its value applies to every pool and takes
precedence over the per-pool ``include_deferred`` column. Returns None when unset.
"""
from airflow.configuration import conf

if conf.get("core", "pool_include_deferred", fallback=""):
return conf.getboolean("core", "pool_include_deferred")
return None

@property
def effective_include_deferred(self) -> bool:
"""The ``include_deferred`` value in effect: the cluster-wide config value when fixed, else the pool's own."""
override = Pool.get_include_deferred_override()
return self.include_deferred if override is None else override

@staticmethod
@provide_session
def get_pools(*, session: Session = NEW_SESSION) -> Sequence[Pool]:
Expand Down Expand Up @@ -143,6 +163,10 @@ def create_or_update_pool(
"team_name cannot be set when multi_team mode is disabled. Please contact your administrator."
)

include_deferred_override = Pool.get_include_deferred_override()
if include_deferred_override is not None:
include_deferred = include_deferred_override

pool = session.scalar(select(Pool).filter_by(pool=name))
if pool is None:
pool = Pool(
Expand Down Expand Up @@ -197,6 +221,7 @@ def slots_stats(

pools: dict[str, PoolStats] = {}
pool_includes_deferred: dict[str, bool] = {}
include_deferred_override = Pool.get_include_deferred_override()

# The below type annotation is acceptable on SQLA2.1, but not on 2.0
query: Select[str, int, bool] = select(Pool.pool, Pool.slots, Pool.include_deferred) # type: ignore[type-arg]
Expand All @@ -210,7 +235,9 @@ def slots_stats(
pools[pool_name] = PoolStats(
total=total_slots, running=0, queued=0, open=0, deferred=0, scheduled=0
)
pool_includes_deferred[pool_name] = include_deferred
pool_includes_deferred[pool_name] = (
include_deferred if include_deferred_override is None else include_deferred_override
)

allowed_execution_states = EXECUTION_STATES | {
TaskInstanceState.DEFERRED,
Expand Down Expand Up @@ -287,7 +314,7 @@ def occupied_slots(self, *, session: Session = NEW_SESSION) -> int:
)

def get_occupied_states(self):
if self.include_deferred:
if self.effective_include_deferred:
return EXECUTION_STATES | {
TaskInstanceState.DEFERRED,
}
Expand Down
11 changes: 11 additions & 0 deletions airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8928,6 +8928,17 @@ export const $ConfigResponse = {
}
],
title: 'Rerun With Latest Version'
},
pool_include_deferred: {
anyOf: [
{
type: 'boolean'
},
{
type: 'null'
}
],
title: 'Pool Include Deferred'
}
},
type: 'object',
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2266,6 +2266,7 @@ export type ConfigResponse = {
theme: Theme | null;
multi_team: boolean;
rerun_with_latest_version?: boolean | null;
pool_include_deferred?: boolean | null;
};

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,7 @@
"checkbox": "Check to include deferred tasks when calculating open pool slots",
"description": "Description",
"includeDeferred": "Include Deferred",
"includeDeferredFixedHelperText": "This option is fixed to \"{{value}}\" for all pools by the cluster-level configuration and cannot be changed per pool.",
"nameMaxLength": "Name can contain a maximum of 256 characters",
"nameRequired": "Name is required",
"slots": "Slots",
Expand Down
23 changes: 20 additions & 3 deletions airflow-core/src/airflow/ui/src/pages/Pools/PoolForm.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -56,9 +56,15 @@ const PoolForm = ({ error, initialPool, isPending, manageMutate, setError }: Poo
mode: "onChange",
});
const multiTeamEnabled = Boolean(useConfig("multi_team"));
const includeDeferredConfig = useConfig("pool_include_deferred");
// A boolean means include_deferred is fixed cluster-wide and cannot be chosen per pool
const includeDeferredOverride =
typeof includeDeferredConfig === "boolean" ? includeDeferredConfig : undefined;

const onSubmit = (data: PoolBody) => {
manageMutate(data);
manageMutate(
includeDeferredOverride === undefined ? data : { ...data, include_deferred: includeDeferredOverride },
);
};

const handleReset = () => {
Expand Down Expand Up @@ -141,11 +147,22 @@ const PoolForm = ({ error, initialPool, isPending, manageMutate, setError }: Poo
control={control}
name="include_deferred"
render={({ field }) => (
<Field.Root mb={4} mt={4}>
<Field.Root disabled={includeDeferredOverride !== undefined} mb={4} mt={4}>
<Field.Label fontSize="md">{translate("pools.form.includeDeferred")}</Field.Label>
<Checkbox checked={field.value} onChange={field.onChange}>
<Checkbox
checked={includeDeferredOverride ?? field.value}
disabled={includeDeferredOverride !== undefined}
onChange={field.onChange}
>
{translate("pools.form.checkbox")}
</Checkbox>
{includeDeferredOverride === undefined ? undefined : (
<Field.HelperText>
{translate("pools.form.includeDeferredFixedHelperText", {
value: includeDeferredOverride ? "True" : "False",
})}
</Field.HelperText>
)}
</Field.Root>
)}
/>
Expand Down
2 changes: 1 addition & 1 deletion airflow-core/src/airflow/utils/db.py
Original file line number Diff line number Diff line change
Expand Up @@ -191,7 +191,7 @@ def add_default_pool_if_not_exists(*, session: Session = NEW_SESSION):
pool=Pool.DEFAULT_POOL_NAME,
slots=conf.getint(section="core", key="default_pool_task_slot_count"),
description="Default pool",
include_deferred=False,
include_deferred=Pool.get_include_deferred_override() or False,
)
session.add(default_pool)
session.commit()
Expand Down
Loading
Loading