Skip to content

Commit 60390d3

Browse files
authored
py: wait for a failed pipeline to shutdown (feldera#2435)
Wait for 15 seconds for a pipeline to shutdown. If the pipeline doesn't shutdown even after the timeout, resend the shutdown request, and wait a little longer. Signed-off-by: Abhinav Gyawali <22275402+abhizer@users.noreply.github.com>
1 parent 4e31f15 commit 60390d3

4 files changed

Lines changed: 58 additions & 6 deletions

File tree

python/feldera/enums.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -181,3 +181,6 @@ def from_str(value):
181181
if member.name.lower() == value.lower():
182182
return member
183183
raise ValueError(f"Unknown value '{value}' for enum {PipelineStatus.__name__}")
184+
185+
def __eq__(self, other):
186+
return self.value == other.value

python/feldera/rest/feldera_client.py

Lines changed: 25 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -225,18 +225,38 @@ def shutdown_pipeline(self, pipeline_name: str):
225225
path=f"/pipelines/{pipeline_name}/shutdown",
226226
)
227227

228-
while True:
228+
start = time.time()
229+
timeout = 15
230+
231+
while time.time() - start < timeout:
229232
status = self.get_pipeline(pipeline_name).deployment_status
230233

231234
if status == "Shutdown":
232-
break
233-
elif status == "Failed":
234-
raise RuntimeError(f"Failed to shutdown pipeline")
235+
return
235236

236237
logging.debug("still shutting down %s, waiting for 100 more milliseconds", pipeline_name)
237238
time.sleep(0.1)
238239

239-
# TODO: better name for this method
240+
# retry sending shutdown request as the pipline hasn't shutdown yet
241+
logging.debug("pipeline %s hasn't shutdown after %s s, retrying", pipeline_name, timeout)
242+
self.http.post(
243+
path=f"/pipelines/{pipeline_name}/shutdown",
244+
)
245+
246+
start = time.time()
247+
timeout = 5
248+
249+
while time.time() - start < timeout:
250+
status = self.get_pipeline(pipeline_name).deployment_status
251+
252+
if status == "Shutdown":
253+
return
254+
255+
logging.debug("still shutting down %s, waiting for 100 more milliseconds", pipeline_name)
256+
time.sleep(0.1)
257+
258+
raise RuntimeError(f"Failed to shutdown pipeline {pipeline_name}")
259+
240260
def push_to_pipeline(
241261
self,
242262
pipeline_name: str,

python/tests/requirements.txt

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,2 +1,2 @@
1-
kafka-python==2.0.2
1+
kafka-python-ng==2.2.2
22
pytest

python/tests/test_pipeline_builder.py

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,11 @@
11
import os
2+
import time
23
import unittest
34
import pandas as pd
45
from kafka import KafkaProducer, KafkaConsumer
56
from kafka.admin import KafkaAdminClient, NewTopic
67

8+
from build.lib.feldera.enums import PipelineStatus
79
from feldera import PipelineBuilder, Pipeline
810
from tests import TEST_CLIENT
911

@@ -699,6 +701,7 @@ def test_pandas_binary(self):
699701
got = out.to_dict()
700702

701703
assert expected_data == got
704+
pipeline.delete()
702705

703706
def test_pandas_decimal(self):
704707
from decimal import Decimal
@@ -721,6 +724,7 @@ def test_pandas_decimal(self):
721724
got = out.to_dict()
722725

723726
assert expected == got
727+
pipeline.delete()
724728

725729
def test_pandas_array(self):
726730
sql = f"""
@@ -742,6 +746,7 @@ def test_pandas_array(self):
742746
datum.update({"insert_delete": 1})
743747

744748
assert got == data
749+
pipeline.delete()
745750

746751
def test_pandas_struct(self):
747752
sql = f"""
@@ -765,6 +770,7 @@ def test_pandas_struct(self):
765770
datum.update({"insert_delete": 1})
766771

767772
assert data == got
773+
pipeline.delete()
768774

769775
def test_pandas_date_time_timestamp(self):
770776
from pandas import Timestamp, Timedelta
@@ -785,6 +791,7 @@ def test_pandas_date_time_timestamp(self):
785791
got = out.to_dict()
786792

787793
assert expected == got
794+
pipeline.delete()
788795

789796
def test_pandas_simple(self):
790797
sql = f"""
@@ -805,6 +812,7 @@ def test_pandas_simple(self):
805812
datum.update({"insert_delete": 1})
806813

807814
assert data == got
815+
pipeline.delete()
808816

809817
def test_pandas_map(self):
810818
sql = f"""
@@ -823,6 +831,27 @@ def test_pandas_map(self):
823831
got = out.to_dict()
824832

825833
assert expected == got
834+
pipeline.delete()
835+
836+
def test_failed_pipeline_shutdown(self):
837+
sql = f"""
838+
CREATE TABLE t0 (c1 TINYINT);
839+
CREATE VIEW v0 AS SELECT c1 + 127::TINYINT FROM t0;"""
840+
841+
pipeline = PipelineBuilder(TEST_CLIENT, name="test_failed_pipeline_shutdown", sql=sql).create_or_replace()
842+
pipeline.start()
843+
data = [{"c1": 127}]
844+
pipeline.input_json("t0", data)
845+
846+
while True:
847+
status = pipeline.status()
848+
expected = PipelineStatus.FAILED
849+
if status == expected:
850+
break
851+
time.sleep(1)
852+
853+
pipeline.shutdown()
854+
pipeline.delete()
826855

827856

828857
if __name__ == '__main__':

0 commit comments

Comments
 (0)