Skip to content

Required to fix error handling in EksCreateClusterOperator in deferrable mode #34844

Description

@Taragolis

Body

There is couple issues happen if something went wrong during cluster creation.

We try to delete cluster by defer triggerer EksDeleteClusterTrigger to method execute_failed

elif event["status"] == "failed":
self.log.error("Cluster failed to start and will be torn down.")
self.eks_hook.delete_cluster(name=self.cluster_name)
self.defer(
trigger=EksDeleteClusterTrigger(
cluster_name=self.cluster_name,
waiter_delay=self.waiter_delay,
waiter_max_attempts=self.waiter_max_attempts,
aws_conn_id=self.aws_conn_id,
region_name=self.region,
force_delete_compute=False,
),
method_name="execute_failed",
timeout=timedelta(seconds=self.waiter_max_attempts * self.waiter_delay),
)

However execute_failed method tried to validate incorrect status (typo in delteted)

def execute_failed(self, context: Context, event: dict[str, Any] | None = None) -> None:
if event is None:
self.log.info("Trigger error: event is None")
raise AirflowException("Trigger error: event is None")
elif event["status"] == "delteted":
self.log.info("Cluster deleted")
raise event["exception"]

This triggerer only returns status of execution and not exception, I guess we can't serialise exception object.

yield TriggerEvent({"status": "deleted"})

Required to fix current error handling in EksDeleteClusterTrigger in deferrable mode. Seems like it could return success status even if execution failed, so I guess we need always raise an error in execute_failed method

Committer

  • I acknowledge that I am a maintainer/committer of the Apache Airflow project.

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions