Skip to content

Commit a65f302

Browse files
committed
edited the comment to add more precise description
1 parent bdde697 commit a65f302

File tree

1 file changed

+4
-3
lines changed

1 file changed

+4
-3
lines changed

python/pyspark/streaming_tests.py

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -50,8 +50,8 @@ def setUp(self):
5050
self.ssc = StreamingContext(appName=class_name, duration=Seconds(1))
5151

5252
def tearDown(self):
53-
# Do not call StreamingContext.stop directly because we do not wait to shutdown
54-
# call back server and py4j client
53+
# Do not call pyspark.streaming.context.StreamingContext.stop directly because
54+
# we do not wait to shutdowncall back server and py4j client
5555
self.ssc._jssc.stop()
5656
self.ssc._sc.stop()
5757
# Why does it long time to terminaete StremaingContext and SparkContext?
@@ -146,7 +146,7 @@ def _run_stream(self, test_input, test_func, expected_output):
146146
"""Start stream and return the output"""
147147
# Generate input stream with user-defined input
148148
test_input_stream = self.ssc._testInputStream(test_input)
149-
# Applied test function to stream
149+
# Apply test function to stream
150150
test_stream = test_func(test_input_stream)
151151
# Add job to get output from stream
152152
test_stream._test_output(StreamOutput.result)
@@ -160,6 +160,7 @@ def _run_stream(self, test_input, test_func, expected_output):
160160
if (current_time - start_time) > self.timeout:
161161
break
162162
self.ssc.awaitTermination(50)
163+
# check if the output is the same length of expexted output
163164
if len(expected_output) == len(StreamOutput.result):
164165
break
165166
return StreamOutput.result

0 commit comments

Comments
 (0)