Skip to content

Allow @task.kubernetes to derive pod name based on decorated function (or more flexibility to override the behaviour) #44779

Description

@wasabigeek

Description

I'd like for @task.kubernetes to derive the task_id / pod_name based on the decorated function's name. To achieve this now, I would need to pass it into the decorator args, which is a bit repetitive:

@task.kubernetes(task_id="some_func", name=f"{some_prefix}_some_func")
def some_func():
    pass

We managed to do something like the above but it was more involved than expected:

# inheriting from a private class, not ideal
class DefaultKubernetesDecoratedOperator(_KubernetesDecoratedOperator):
    custom_operator_name = "@default_kubernetes_task" # the operator removes the decorator so this is required

    def __init__(self, **kwargs: Any) -> None:
        if "task_id" not in kwargs and "python_callable" in kwargs:
            task_id = kwargs["python_callable"].__name__
            kwargs["task_id"] = task_id

        if "task_id" in kwargs and "name" not in kwargs:
            name = f"{get_prefix()}-{kwargs['task_id']}"
            kwargs["name"] = name

        super().__init__(**kwargs)


def default_kubernetes_task(
    python_callable: Callable[[], None] | None = None,
    multiple_outputs: bool | None = None,
    **kwargs: Any,
) -> TaskDecorator:
    return task_decorator_factory(
        python_callable=python_callable,
        multiple_outputs=multiple_outputs,
        decorated_operator_class=DefaultKubernetesDecoratedOperator,
        **kwargs,
    )

Initially I tried wrapping task.kubernetes itself, but the source code scrubbing relies on a hardcoded decorator name, so it doesn't scrub the decorator and the pod fails with NameError: name 'default_kubernetes_task' is not defined:

def default_kubernetes_task(
    **task_kwargs: Any,
) -> Callable[[Callable[..., Any]], Callable[..., Any]]:
    def decorator(func: Callable[..., Any]) -> Any:
        if "task_id" not in task_kwargs:
            task_kwargs["task_id"] = func.__name__

        if "name" not in task_kwargs:
            name = f"{get_[prefix()}-{task_kwargs['task_id']}"
            task_kwargs["name"] = name

        return task.kubernetes(**task_kwargs)(func)

    return decorator

It'd be nice if there were a more official way to achieve the above (or something similar). Perhaps:

  • @task.kubernetes could accept an override for the decorated_operator_class?
  • a list of custom operator names to remove could be set in the configuration?

Use case/motivation

I'd like for @task.kubernetes to derive the task_id / pod_name based on the decorated function's name, but more generally it could be useful to allow task.kubernetes to be overridden?

Related issues

No response

Are you willing to submit a PR?

  • Yes I am willing to submit a PR!

Code of Conduct

Activity

  1. boring-cyborg commented on Dec 9, 2024

    @boring-cyborg

    Thanks for opening your first issue here! Be sure to follow the issue template! If you are willing to raise PR to address this issue please do so, no need to wait for approval.

  2. kyleburke-meq commented on Jan 27, 2025

    @kyleburke-meq

    I would also like this to be more flexible. Thank you for pointing out that i need to override the custom_operator_name so it can strip the decorator name.

  3. removed
    needs-triagelabel for new issues that we didn't triage yet
    on Feb 9, 2025
  4. insomnes commented on Feb 28, 2025

    @insomnes
    Contributor

    I believe the pod naming part was covered in:

    #46464

    I don't get it with the task id: isn't it already defaults to the decorated python callable?

  5. RNHTTR commented on Mar 12, 2025

    @RNHTTR
    Contributor

    I agree with @insomnes -- I believe this has been resolved via PR #46535. Does version 10.3.0 of apache-airflow-providers-cncf-kubernetes solve your problem?

  6. github-actions commented on Mar 27, 2025

    @github-actions
    Contributor

    This issue has been automatically marked as stale because it has been open for 14 days with no response from the author. It will be closed in next 7 days if no further activity occurs from the issue author.

  7. added
    staleStale PRs per the .github/workflows/stale.yml policy file
    on Mar 27, 2025
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions