Repository navigation
Conversation
One potential problem we have with batch processing is that any one slow item will clog up the whole batch. This pr implements a queueing method instead, where we keep N queues that each have their own workers. There's still a chance of individual items backlogging a queue, but we can try increased concurrency here to reduce the chances of that happening <!-- Describe your PR here. -->
Comment on lines
+49
to
+54
| def _get_partition_lock(self, partition: Partition) -> threading.Lock: | ||
| """Get or create a lock for a partition.""" | ||
| lock = self.partition_locks.get(partition) | ||
| if lock: | ||
| return lock | ||
| return self.partition_locks.setdefault(partition, threading.Lock()) |
There was a problem hiding this comment.
There's a potential thread safety issue in the _get_partition_lock() method. The current implementation has a race condition between checking if the lock exists and creating it with setdefault(). If two threads simultaneously determine that a lock doesn't exist for a partition, both could attempt to create one, resulting in different lock objects for the same partition.
Consider replacing lines 51-54 with a single atomic operation:
return self.partition_locks.setdefault(partition, threading.Lock())This ensures that only one lock is ever created per partition, maintaining proper synchronization across threads.
Suggested change
| def _get_partition_lock(self, partition: Partition) -> threading.Lock: | |
| """Get or create a lock for a partition.""" | |
| lock = self.partition_locks.get(partition) | |
| if lock: | |
| return lock | |
| return self.partition_locks.setdefault(partition, threading.Lock()) | |
| def _get_partition_lock(self, partition: Partition) -> threading.Lock: | |
| """Get or create a lock for a partition.""" | |
| return self.partition_locks.setdefault(partition, threading.Lock()) | |
Spotted by Diamond
Is this helpful? React 👍 or 👎 to let us know.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Test 9