Skip to content

Commit b2d4901

Browse files
committed
tweak
1 parent 47d08a7 commit b2d4901

1 file changed

Lines changed: 3 additions & 1 deletion

File tree

streamz/tests/test_kafka.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -169,7 +169,7 @@ def test_from_kafka_thread():
169169
stream = Stream.from_kafka([TOPIC], ARGS)
170170
out = stream.sink_to_list()
171171
stream.start()
172-
yield gen.sleep(1.1)
172+
yield await_for(stream.started, 10, period=0.1)
173173
for i in range(10):
174174
yield gen.sleep(0.1)
175175
kafka.produce(TOPIC, b'value-%d' % i)
@@ -585,6 +585,8 @@ def test_kafka_checkpointing_auto_offset_reset_latest():
585585

586586
stream1 = Stream.from_kafka_batched(TOPIC, ARGS, asynchronous=True)
587587
out1 = stream1.map(split).gather().sink_to_list()
588+
time.sleep(1) # messages make ttheir way through kafka
589+
588590
stream1.start()
589591
wait_for(lambda: stream1.upstream.started, 10, period=0.1)
590592

0 commit comments

Comments
 (0)