We read every piece of feedback, and take your input very seriously.
To see all available qualifiers, see our documentation.
1 parent 47d08a7 commit b2d4901Copy full SHA for b2d4901
1 file changed
streamz/tests/test_kafka.py
@@ -169,7 +169,7 @@ def test_from_kafka_thread():
169
stream = Stream.from_kafka([TOPIC], ARGS)
170
out = stream.sink_to_list()
171
stream.start()
172
- yield gen.sleep(1.1)
+ yield await_for(stream.started, 10, period=0.1)
173
for i in range(10):
174
yield gen.sleep(0.1)
175
kafka.produce(TOPIC, b'value-%d' % i)
@@ -585,6 +585,8 @@ def test_kafka_checkpointing_auto_offset_reset_latest():
585
586
stream1 = Stream.from_kafka_batched(TOPIC, ARGS, asynchronous=True)
587
out1 = stream1.map(split).gather().sink_to_list()
588
+ time.sleep(1) # messages make ttheir way through kafka
589
+
590
stream1.start()
591
wait_for(lambda: stream1.upstream.started, 10, period=0.1)
592
0 commit comments