Skip to content

Commit

Permalink
fix: add lock to flow control (#899)
Browse files Browse the repository at this point in the history
  • Loading branch information
daniel-sanche authored Dec 12, 2023
1 parent fe58f61 commit e4e63c7
Showing 1 changed file with 7 additions and 4 deletions.
11 changes: 7 additions & 4 deletions google/cloud/bigtable/batcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ def __init__(
self.inflight_size = 0
self.event = threading.Event()
self.event.set()
self._lock = threading.Lock()

def is_blocked(self):
"""Returns True if:
Expand All @@ -132,8 +133,9 @@ def control_flow(self, batch_info):
Calculate the resources used by this batch
"""

self.inflight_mutations += batch_info.mutations_count
self.inflight_size += batch_info.mutations_size
with self._lock:
self.inflight_mutations += batch_info.mutations_count
self.inflight_size += batch_info.mutations_size
self.set_flow_control_status()

def wait(self):
Expand All @@ -158,8 +160,9 @@ def release(self, batch_info):
Release the resources.
Decrement the row size to allow enqueued mutations to be run.
"""
self.inflight_mutations -= batch_info.mutations_count
self.inflight_size -= batch_info.mutations_size
with self._lock:
self.inflight_mutations -= batch_info.mutations_count
self.inflight_size -= batch_info.mutations_size
self.set_flow_control_status()


Expand Down

0 comments on commit e4e63c7

Please sign in to comment.