From 12ddb1d09e6d2039f018a355954e31c0ae62f372 Mon Sep 17 00:00:00 2001 From: Khor Shu Heng Date: Tue, 3 Nov 2020 15:01:48 +0800 Subject: [PATCH] Wait for job to be ready before cancelling Signed-off-by: Khor Shu Heng --- tests/integration/test_launchers.py | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/tests/integration/test_launchers.py b/tests/integration/test_launchers.py index fb6731afa2d..4b80afe054d 100644 --- a/tests/integration/test_launchers.py +++ b/tests/integration/test_launchers.py @@ -1,6 +1,6 @@ -from time import sleep +import time -from feast.pyspark.abc import RetrievalJobParameters, SparkJobStatus +from feast.pyspark.abc import RetrievalJobParameters, SparkJobStatus, SparkJob from feast.pyspark.launchers.gcloud import DataprocClusterLauncher from .fixtures.job_parameters import customer_entity # noqa: F401 @@ -9,6 +9,13 @@ from .fixtures.launchers import dataproc_launcher # noqa: F401 +def wait_for_job_status(job: SparkJob, expected_status: SparkJobStatus, max_retry: int = 4, retry_interval: int = 5): + for i in range(max_retry): + if job.get_status() == expected_status: + return + time.sleep(retry_interval) + raise ValueError(f"Timeout waiting for job status to become {expected_status.name}") + def test_dataproc_job_api( dataproc_launcher: DataprocClusterLauncher, # noqa: F811 dataproc_retrieval_job_params: RetrievalJobParameters, # noqa: F811 @@ -27,5 +34,6 @@ def test_dataproc_job_api( job.get_id() for job in dataproc_launcher.list_jobs(include_terminated=False) ] assert job_id in active_job_ids + wait_for_job_status(retrieved_job, SparkJobStatus.IN_PROGRESS) retrieved_job.cancel() assert retrieved_job.get_status() == SparkJobStatus.FAILED