Skip to content

Commit 394428c

Browse files
committed
Move lock from _next_id to register.
1 parent 206a100 commit 394428c

1 file changed

Lines changed: 18 additions & 43 deletions

File tree

sdks/python/apache_beam/utils/subprocess_server.py

Lines changed: 18 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -76,63 +76,38 @@ def __init__(self, constructor, destructor):
7676
self._cache = {}
7777
self._lock = threading.Lock()
7878
self._counter = 0
79-
self._log_buffer = []
79+
80+
def _next_id(self):
81+
# Caller must hold self._lock.
82+
self._counter += 1
83+
return self._counter
8084

8185
def register(self):
82-
import traceback
83-
stack = "".join(traceback.format_stack())
8486
with self._lock:
85-
self._counter += 1
86-
owner = self._counter
87+
owner = self._next_id()
8788
self._live_owners.add(owner)
88-
msg = f"Registered subprocess owner {owner} (PID {os.getpid()})"
89-
self._log_buffer.append((msg, stack))
9089
return owner
9190

9291
def purge(self, owner):
9392
to_delete = []
94-
should_flush = False
95-
import traceback
96-
stack = "".join(traceback.format_stack())
97-
9893
with self._lock:
9994
if owner not in self._live_owners:
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-
95+
_LOGGER.warning(
96+
"Subprocess owner %s already purged. If this occurs during atexit "
97+
"shutdown, the subprocess was already cleaned up earlier.",
98+
owner)
99+
return
100+
self._live_owners.remove(owner)
101+
for key, entry in list(self._cache.items()):
102+
if owner in entry.owners:
103+
entry.owners.remove(owner)
104+
if not entry.owners:
105+
to_delete.append(entry.obj)
106+
del self._cache[key]
124107
# Actually call the destructors outside of the lock.
125108
for value in to_delete:
126109
self._destructor(value)
127110

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-
136111
def get(self, *key):
137112
if not self._live_owners:
138113
raise RuntimeError("At least one owner must be registered.")

0 commit comments

Comments
 (0)