Skip to content

Commit 206a100

Browse files
committed
Store the registration and purge in memory and output at the end.
1 parent d0487f7 commit 206a100

1 file changed

Lines changed: 44 additions & 22 deletions

File tree

sdks/python/apache_beam/utils/subprocess_server.py

Lines changed: 44 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -76,41 +76,63 @@ def __init__(self, constructor, destructor):
7676
self._cache = {}
7777
self._lock = threading.Lock()
7878
self._counter = 0
79+
self._log_buffer = []
7980

80-
def _next_id(self):
81+
def register(self):
82+
import traceback
83+
stack = "".join(traceback.format_stack())
8184
with self._lock:
8285
self._counter += 1
83-
return self._counter
84-
85-
def register(self):
86-
owner = self._next_id()
87-
self._live_owners.add(owner)
88-
_LOGGER.warning("Registered subprocess owner %s (PID %s)", owner, os.getpid(), stack_info=True)
86+
owner = self._counter
87+
self._live_owners.add(owner)
88+
msg = f"Registered subprocess owner {owner} (PID {os.getpid()})"
89+
self._log_buffer.append((msg, stack))
8990
return owner
9091

9192
def purge(self, owner):
9293
to_delete = []
94+
should_flush = False
95+
import traceback
96+
stack = "".join(traceback.format_stack())
97+
9398
with self._lock:
9499
if owner not in self._live_owners:
95-
_LOGGER.warning(
96-
"Subprocess owner %s (PID %s) already purged. If this occurs during atexit "
97-
"shutdown, the subprocess was already cleaned up earlier.",
98-
owner,
99-
os.getpid(),
100-
stack_info=True)
101-
return
102-
self._live_owners.remove(owner)
103-
_LOGGER.warning("Purging subprocess owner %s (PID %s)", owner, os.getpid(), stack_info=True)
104-
for key, entry in list(self._cache.items()):
105-
if owner in entry.owners:
106-
entry.owners.remove(owner)
107-
if not entry.owners:
108-
to_delete.append(entry.obj)
109-
del self._cache[key]
100+
msg = (
101+
f"Subprocess owner {owner} (PID {os.getpid()}) already purged. "
102+
"If this occurs during atexit shutdown, the subprocess was already cleaned up earlier."
103+
)
104+
self._log_buffer.append((msg, stack))
105+
should_flush = True
106+
else:
107+
self._live_owners.remove(owner)
108+
msg = f"Purging subprocess owner {owner} (PID {os.getpid()})"
109+
self._log_buffer.append((msg, stack))
110+
111+
for key, entry in list(self._cache.items()):
112+
if owner in entry.owners:
113+
entry.owners.remove(owner)
114+
if not entry.owners:
115+
to_delete.append(entry.obj)
116+
del self._cache[key]
117+
118+
if to_delete:
119+
should_flush = True
120+
121+
if should_flush:
122+
self._flush_logs()
123+
110124
# Actually call the destructors outside of the lock.
111125
for value in to_delete:
112126
self._destructor(value)
113127

128+
def _flush_logs(self):
129+
with self._lock:
130+
logs_to_print = list(self._log_buffer)
131+
self._log_buffer.clear()
132+
133+
for msg, stack in logs_to_print:
134+
_LOGGER.warning("%s\nStack Trace:\n%s", msg, stack)
135+
114136
def get(self, *key):
115137
if not self._live_owners:
116138
raise RuntimeError("At least one owner must be registered.")

0 commit comments

Comments
 (0)