Skip to content

Possible deadlock in SitemapRequestLoader #1842

Description

@bdm83

I've been trying to debug stuck crawlee process that uses the SitemapRequestLoader. The crawler will correctly crawl for several hours (12+) but will eventually stop processing new requests and enter a blocked state.

I was able to perform an asyncio task dump while the program was in a stuck state and noticed these two tasks:

--- Task 16: Task-2 ---
  Done: False
  Cancelled: False
  Coro: <coroutine object SitemapRequestLoader._load_sitemaps at 0x7f9f67f7d000>
    File ".../crawlee/request_loaders/_sitemap_request_loader.py", line 261, in _load_sitemaps
    await self._queue_has_capacity.wait()

--- Task 6: autoscaled pool worker task orchestrator ---
  Done: False
  Cancelled: False
  Coro: <coroutine object AutoscaledPool._worker_task_orchestrator at 0x7f9f6751b340>
    File ".../crawlee/_autoscaling/autoscaled_pool.py", line 245, in _worker_task_orchestrator
    await asyncio.wait_for(run.worker_tasks_updated.wait(), timeout=0.5)

Task 16 (producer) is stuck on _queue_has_capacity.wait(). Task 6 (orchestrator) is idle-polling because the queue is empty and no workers are running

From inspecting _sitemap_request_loader.py it looks like there could be a logic bug with setting/clearing self._queue_has_capacity:

async def _load_sitemaps(self) -> None:
    ...
    # Check if we have capacity in the queue
    await self._queue_has_capacity.wait()
                                           # Possible bug if consumer
    state = await self._get_state()        # sets _queue_has_capacity while
    async with self._queue_lock:           # producer is here
        state.url_queue.append(url)
        state.current_sitemap_processed_urls.add(url)
        state.total_count += 1
        if len(state.url_queue) >= self._max_buffer_size:
            # Notify that the queue is full
            self._queue_has_capacity.clear()
    ...
@override
async def fetch_next_request(self) -> Request | None:
    ...
    async with self._queue_lock:
        url = state.url_queue.popleft()
        ...
        if len(state.url_queue) < self._max_buffer_size:
            self._queue_has_capacity.set()
    ...

From my (limited) understanding of how this should run, it seems like if fetch_next_request(self) is called by a consumer JUST BEFORE the producer starts waiting on self._queue_lock (_load_sitemaps) then the self._queue_has_capacity.set() generated by the consumer will have no effect (since it has already been set) and the next time await self._queue_has_capacity.wait() is called, it will be stuck. This is due to the self._queue_has_capacity.set() generated by the consumer being "wasted" because it is called BEFORE self._queue_has_capacity gets cleared.

Activity

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

Metadata

Metadata

Assignees

Labels

bugSomething isn't working.t-toolingIssues with this label are in the ownership of the tooling team.

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions