@@ -360,23 +360,31 @@ def test_batch_write_with_ordering_key(self):
360360
361361 time .sleep (10 )
362362
363- response = self .sub_client .pull (
364- request = {
365- 'subscription' : ordering_sub .name ,
366- 'max_messages' : 10 ,
367- })
368-
369- self .assertEqual (len (response .received_messages ), len (test_messages ))
363+ # Retry pulling to handle PubSub delivery delays
364+ received_messages = []
365+ deadline = time .time () + 60 # wait up to 60 seconds
366+ while time .time () < deadline :
367+ response = self .sub_client .pull (
368+ request = {
369+ 'subscription' : ordering_sub .name ,
370+ 'max_messages' : 10 ,
371+ })
372+ received_messages .extend (response .received_messages )
373+ if len (received_messages ) >= len (test_messages ):
374+ break
375+ time .sleep (5 )
376+
377+ self .assertEqual (len (received_messages ), len (test_messages ))
370378
371379 received_map = {
372380 msg .message .data : msg .message
373- for msg in response . received_messages
381+ for msg in received_messages
374382 }
375383 self .assertEqual (received_map [b'order_data001' ].ordering_key , 'key1' )
376384 self .assertEqual (received_map [b'order_data002' ].ordering_key , 'key1' )
377385 self .assertEqual (received_map [b'order_data003' ].ordering_key , 'key2' )
378386
379- ack_ids = [msg .ack_id for msg in response . received_messages ]
387+ ack_ids = [msg .ack_id for msg in received_messages ]
380388 self .sub_client .acknowledge (
381389 request = {
382390 'subscription' : ordering_sub .name ,
0 commit comments