Skip to content

feat(uptime): Add ability to use queues to manage parallelism - #9

Open
everettbu wants to merge 1 commit into
kafka-consumer-parallel-beforefrom
kafka-consumer-parallel-after
Open

everettbu wants to merge 1 commit into
kafka-consumer-parallel-beforefrom
kafka-consumer-parallel-after

Conversation

@everettbu

@everettbu everettbu commented Jul 26, 2025 •

Copy link
Copy Markdown

Test 9

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())

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants