@@ -474,8 +474,8 @@ def __init__(self, data_buffer_time_limit_ms=0):
474474 self ._start_send_forwarder ()
475475 self ._received = collections .defaultdict (
476476 lambda : ByteLimitedQueue (
477- maxsize = _DEFAULT_RECEIVE_QUEUE_MAX_ELEMENTS ,
478- maxbytes = _DEFAULT_SEND_QUEUE_MAX_BYTES )
477+ maxsize = _DEFAULT_RECEIVE_QUEUE_MAX_ELEMENTS , maxbytes =
478+ _DEFAULT_SEND_QUEUE_MAX_BYTES )
479479 ) # type: DefaultDict[str, ByteLimitedQueue[DataOrTimers]]
480480
481481 # Keep a cache of completed instructions. Data for completed instructions
@@ -494,27 +494,29 @@ def close(self):
494494 self ._enqueue_to_send (self ._WRITES_FINISHED )
495495 if self ._send_forwarder is not None :
496496 self ._send_forwarder .join ()
497+ if self ._exception :
498+ raise self ._exception
497499 self ._closed = True
498500
499501 def _start_send_forwarder (self ):
500502 # type: () -> None
501503 forwarder = threading .Thread (
502- target = self ._forward_pending_to_send ,
503- name = 'forward_grpc_outputs' )
504+ target = self ._forward_pending_to_send , name = 'forward_grpc_outputs' )
504505 forwarder .daemon = True
505506 forwarder .start ()
506507 self ._send_forwarder = forwarder
507508
508509 def _enqueue_to_send (self , elem ):
509510 # type: (DataOrTimers) -> None
510- self ._pending_send .put (elem , self ._get_element_size_bytes (elem ))
511+ size = self ._get_element_size_bytes (elem )
512+ self ._pending_send .put ((elem , size ), size )
511513
512514 def _forward_pending_to_send (self ):
513515 # type: () -> None
514516 try :
515517 while True :
516- elem = self ._pending_send .get ()
517- self ._to_send .put (elem , self . _get_element_size_bytes ( elem ) )
518+ elem , size = self ._pending_send .get ()
519+ self ._to_send .put (elem , size )
518520 if elem is self ._WRITES_FINISHED :
519521 return
520522 except Exception as e :
0 commit comments