I am listening for Google PubSub messages using the Python code at the bottom of this question. It's effectively the asynchronous pull example from Google.
I run my program and output to file:
python my_script.py | tee log.txt
If I run the program whilst messages are being received, the two print() statements output and everything works as expected.
However, if I run the program before messages are published there is no output. The two print() statements do not output, the code just blocks. I expected to see the two print() statements.
I tried python my_script.py > log.txt but it made no difference.
If I Ctrl + C the program the stack trace shows this:
^CTraceback (most recent call last):
File "my_script.py", line 58, in <module>
streaming_pull_future.result()
File "python3-3.8.6-env/lib/python3.8/site-packages/google/cloud/pubsub_v1/futures.py", line 102, in result
err = self.exception(timeout=timeout)
File "python3-3.8.6-env/lib/python3.8/site-packages/google/cloud/pubsub_v1/futures.py", line 121, in exception
if not self._completed.wait(timeout=timeout):
File "python3-3.8.6/lib/python3.8/threading.py", line 558, in wait
signaled = self._cond.wait(timeout)
File "python3-3.8.6/lib/python3.8/threading.py", line 302, in wait
waiter.acquire()
KeyboardInterrupt
Exception ignored in: <_io.TextIOWrapper name='<stdout>' mode='w' encoding='utf-8'>
BrokenPipeError: [Errno 32] Broken pipe
My guess is the blocking call to result() is preventing flushing the log output to file?
Is it possible to run the program and see logging regardless of whether messages are currently being published?
If not would the logging get "flushed" eventually when messages are finally published?
Code:
subscriber = pubsub_v1.SubscriberClient()
subscription_path = subscriber.subscription_path(project_id, subscription_id )
print(f"Subscribing to {subscription_path}..\n")
streaming_pull_future = subscriber.subscribe(subscription_path, callback=callback)
print(f"Listening for messages on {subscription_path}..\n")
with subscriber:
try:
streaming_pull_future.result()
except TimeoutError:
streaming_pull_future.cancel()