Repository navigation
Publisher thread terminates, forever breaking publication when GCE metadata service blips #1173
Description
Activity
- addedapi: pubsubIssues related to the googleapis/python-pubsub API.Issues related to the googleapis/python-pubsub API.
on May 13, 2024 Hi pgcamus, thank you for bringing this to our attention. Wanted to list some of my observations:
the commit thread is dead forever.
This is correct. But, one thing I observe here is that the publishing of subsequent messages aren't affected, because of this sequence of events:
- TransportError thrown by gapic_publish:
response = self._client._gapic_publish( - Not caught by the Except clause that only catches GoogleAPIError:
except google.api_core.exceptions.GoogleAPIError as exc: - Thread in which
_commitruns dies - Status of the batch instance still remains IN_PROGRESS: , i.e it is not set to
self._status = base.BatchStatus.IN_PROGRESS BatchStatus.ERRORlike it is in case theGoogleAPIErroris thrown:self._status = base.BatchStatus.ERROR
(this causes a separate issue - more on this later)
-
For the next message that is attempted to be published in this batch, it batch.publish() would return a
Nonefuture because of this condition check:python-pubsub/google/cloud/pubsub_v1/publisher/_batch/thread.py
Lines 449 to 450 in cdaf6e9
if self.status != base.BatchStatus.ACCEPTING_MESSAGES: return None -
This would result in the creation of a new batch instance altogether in the sequencer that is trying to publish to the batch:
python-pubsub/google/cloud/pubsub_v1/publisher/_sequencer/unordered_sequencer.py
Lines 156 to 162 in cdaf6e9
future = batch.publish(wrapper) # batch is full, triggering commit_when_full if future is None: batch = self._create_batch(commit_retry=retry, commit_timeout=timeout) # At this point, we lose track of the old batch, but we don't # care since it's already committed (because it was full.) self._current_batch = batch -
This new batch instance would create a new commit thread altogether, whenever the criteria to flush the new batch of messages is met:
self._start_commit_thread()
Having mentioned the above, I do observe that there are a few related issues that need to be addressed:
- Since the TransportError is never caught, and the code path that sets exceptions on the future objects corresponding to messages in the batch that failed never runs. Consequently, futures that are associated with the messages in the batch whose commit thread crashed would never get notified of the exception. They'd indefinitely hang unless they have a timeout set as part of
future.result(timeout=x)(even then, they'd only raise a TimeoutError but not have the exception object set) :python-pubsub/google/cloud/pubsub_v1/publisher/_batch/thread.py
Lines 368 to 369 in cdaf6e9
for future in self._futures: future.set_exception(exc)
Potential fixes include catching that in Batch._commit
This does solve the problem of:
- The commit thread crashing instead of catching the Exception and returning gracefully
- The futures of the messages in the batch that failed having
future.set_exceptiondone on them
The side effect of the above change, however, would be that, given the state of the batch would now be set to
BatchStatus.ERROR, it would cause an AssertionError to be raised for the publish of the subsequent message on the batch:python-pubsub/google/cloud/pubsub_v1/publisher/_batch/thread.py
Lines 445 to 447 in cdaf6e9
assert ( self._status != base.BatchStatus.ERROR ), "Publish after stop() or publish error." which would then be bubbled up until a BaseException is thrown and the client library crashes in the publisher client:
python-pubsub/google/cloud/pubsub_v1/publisher/client.py
Lines 481 to 495 in cdaf6e9
except BaseException as be: # Exceptions can be thrown when attempting to add messages to # the batch. If they're thrown, record them in publisher # batching and create span, end the spans and bubble the # exception up. if self._open_telemetry_enabled: if wrapper: wrapper.end_publisher_batching_span(be) wrapper.end_create_span(be) else: # pragma: NO COVER warnings.warn( message="PublishMessageWrapper is None. Hence, not recording exception and ending publisher batching span and create span", category=RuntimeWarning, ) raise be I'm exploring a few options to solve the above issues, without causing unintended side effects and also factoring in the behavior of the Pub/Sub client libraries of other languages to keep the behavior consistent across the libraries. I'll keep this thread posted. Thanks!
- TransportError thrown by gapic_publish:
Thanks for stopping by to let us know something could be better!
PLEASE READ: If you have a support contract with Google, please create an issue in the support console instead of filing on GitHub. This will ensure a timely response.
Please run down the following list and make sure you've tried the usual "quick fixes":
If you are still having issues, please be sure to include as much information as possible:
Environment details
pip --versiongoogle-cloud-pubsubversion: 2.21.1Steps to reproduce
Run
google-cloud-pubsuband suffer a metadata outage like https://status.cloud.google.com/incidents/u6rQ2nNVbhAFqGCcTm58.Note that this can trigger even in an un-sustained GCE metadata outage as once this exception triggers even once, the commit thread is dead forever. In our case, there was a short outage on the metadata server, but the retries all happened so quickly that the exception was raised before the service recovered
Code example
# exampleStack trace
Speculative analysis
It looks like the issue is that the
google-authlibrary is raising aTransportErrorwhich is not caught by the batch commit thread in this library. Potential fixes include catching that inBatch._commit(e.g. here), or catching it further down ingoogle-cloud-pubsuband wrapping it in aGoogleAPIError.