From 193cc2e14b0256c2bcc735942396159a911a3729 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Thu, 1 Oct 2026 08:19:57 +0000 Subject: [PATCH 01/22] feat: add direct async iteration on Dataset and KeyValueStore Add `__aiter__` to `Dataset` (delegating to `iterate_items`) and to `KeyValueStore` (delegating to the new `iterate_entries`), plus `KeyValueStore.iterate_values` and `KeyValueStore.iterate_entries`. Values are fetched one at a time as the iteration advances, so only a single value is held in memory regardless of the record sizes. On the Apify platform this means one request per record on top of the paginated key listing, which is the only way the API offers to read values. The memory key-value store client now skips keys deleted while a key iteration is suspended instead of raising `KeyError`, matching the other backends. Closes #1745 Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../_memory/_key_value_store_client.py | 6 +- src/crawlee/storages/_dataset.py | 15 ++++ src/crawlee/storages/_key_value_store.py | 59 ++++++++++++++ tests/unit/storages/test_dataset.py | 17 +++++ tests/unit/storages/test_key_value_store.py | 76 +++++++++++++++++++ 5 files changed, 171 insertions(+), 2 deletions(-) diff --git a/src/crawlee/storage_clients/_memory/_key_value_store_client.py b/src/crawlee/storage_clients/_memory/_key_value_store_client.py index e984a9932a..13014099d7 100644 --- a/src/crawlee/storage_clients/_memory/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_memory/_key_value_store_client.py @@ -151,9 +151,11 @@ async def iterate_keys( if limit is not None: keys = keys[:limit] - # Yield metadata for each key + # Yield metadata for each key, skipping records deleted while the iteration was suspended. for key in keys: - record = self._records[key] + record = self._records.get(key) + if record is None: + continue yield KeyValueStoreRecordMetadata( key=key, content_type=record.content_type, diff --git a/src/crawlee/storages/_dataset.py b/src/crawlee/storages/_dataset.py index 1f5a7297ba..0b4aece7ac 100644 --- a/src/crawlee/storages/_dataset.py +++ b/src/crawlee/storages/_dataset.py @@ -246,6 +246,21 @@ async def iterate_items( ): yield item + def __aiter__(self) -> AsyncIterator[Mapping[str, JsonSerializable]]: + """Iterate over all items in the dataset. + + Allows using the dataset directly in an `async for` loop. It is equivalent to calling `iterate_items` + with the default arguments. + + ### Usage + + ```python + async for item in dataset: + print(item) + ``` + """ + return self.iterate_items() + async def list_items( self, *, diff --git a/src/crawlee/storages/_key_value_store.py b/src/crawlee/storages/_key_value_store.py index 557f3432bb..aad9aef674 100644 --- a/src/crawlee/storages/_key_value_store.py +++ b/src/crawlee/storages/_key_value_store.py @@ -213,6 +213,65 @@ async def iterate_keys( ): yield item + async def iterate_values( + self, + exclusive_start_key: str | None = None, + limit: int | None = None, + ) -> AsyncIterator[Any]: + """Iterate over the values of the existing records in the KVS. + + The values are fetched one by one as the iteration advances, so only a single value is held in memory at + a time. On remote backends this means one request per record on top of the paginated key listing. + + Args: + exclusive_start_key: Key to start the iteration from. + limit: Maximum number of records to iterate over. None means no limit. + + Yields: + The value of each record. + """ + async for _, value in self.iterate_entries(exclusive_start_key=exclusive_start_key, limit=limit): + yield value + + async def iterate_entries( + self, + exclusive_start_key: str | None = None, + limit: int | None = None, + ) -> AsyncIterator[tuple[str, Any]]: + """Iterate over the existing records in the KVS as `(key, value)` pairs. + + The values are fetched one by one as the iteration advances, so only a single value is held in memory at + a time. On remote backends this means one request per record on top of the paginated key listing. A record + deleted after its key was listed but before its value was fetched is skipped. + + Args: + exclusive_start_key: Key to start the iteration from. + limit: Maximum number of records to iterate over. None means no limit. + + Yields: + A `(key, value)` tuple for each record. + """ + async for metadata in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit): + record = await self._client.get_value(key=metadata.key) + if record is None: + continue + yield metadata.key, record.value + + def __aiter__(self) -> AsyncIterator[tuple[str, Any]]: + """Iterate over all records in the KVS as `(key, value)` pairs. + + Allows using the key-value store directly in an `async for` loop. It is equivalent to calling + `iterate_entries` with the default arguments. + + ### Usage + + ```python + async for key, value in kvs: + print(key, value) + ``` + """ + return self.iterate_entries() + async def list_keys( self, exclusive_start_key: str | None = None, diff --git a/tests/unit/storages/test_dataset.py b/tests/unit/storages/test_dataset.py index c18c3e54f3..7e0b64fbd5 100644 --- a/tests/unit/storages/test_dataset.py +++ b/tests/unit/storages/test_dataset.py @@ -280,6 +280,23 @@ async def test_iterate_items(dataset: Dataset) -> None: assert collected_items[-1]['id'] == 5 +async def test_async_iteration(dataset: Dataset) -> None: + """Test that the dataset can be used directly in an `async for` loop.""" + items = [{'id': i} for i in range(1, 6)] # 5 items + await dataset.push_data(items) + + collected_items = [item async for item in dataset] + + assert collected_items == items + + +async def test_async_iteration_empty_dataset(dataset: Dataset) -> None: + """Test that iterating over an empty dataset yields nothing.""" + collected_items = [item async for item in dataset] + + assert collected_items == [] + + async def test_iterate_items_with_options(dataset: Dataset) -> None: """Test iterating with offset, limit and desc parameters.""" # Add some items diff --git a/tests/unit/storages/test_key_value_store.py b/tests/unit/storages/test_key_value_store.py index 13b66b4f3a..fc2d4a8f30 100644 --- a/tests/unit/storages/test_key_value_store.py +++ b/tests/unit/storages/test_key_value_store.py @@ -267,6 +267,82 @@ async def test_iterate_keys_with_limit(kvs: KeyValueStore) -> None: assert len(collected_keys) == 5 +async def test_iterate_values(kvs: KeyValueStore) -> None: + """Test iterating over values in the key-value store.""" + await kvs.set_value('key1', 'value1') + await kvs.set_value('key2', {'nested': 2}) + await kvs.set_value('key3', [3]) + + collected_values = [value async for value in kvs.iterate_values()] + + assert len(collected_values) == 3 + assert 'value1' in collected_values + assert {'nested': 2} in collected_values + assert [3] in collected_values + + +async def test_iterate_entries(kvs: KeyValueStore) -> None: + """Test iterating over (key, value) pairs in the key-value store.""" + await kvs.set_value('key1', 'value1') + await kvs.set_value('key2', {'nested': 2}) + await kvs.set_value('key3', [3]) + + collected_entries = dict([entry async for entry in kvs.iterate_entries()]) + + assert collected_entries == {'key1': 'value1', 'key2': {'nested': 2}, 'key3': [3]} + + +async def test_iterate_entries_with_limit_and_exclusive_start_key(kvs: KeyValueStore) -> None: + """Test that `iterate_entries` passes `limit` and `exclusive_start_key` through to the key listing.""" + for i in range(10): + await kvs.set_value(f'key{i}', f'value{i}') + + all_keys = [metadata.key for metadata in await kvs.list_keys()] + start_key = all_keys[2] + + collected_entries = [entry async for entry in kvs.iterate_entries(exclusive_start_key=start_key, limit=3)] + + assert len(collected_entries) == 3 + expected_keys = all_keys[all_keys.index(start_key) + 1 :][:3] + assert [key for key, _ in collected_entries] == expected_keys + assert all(value == f'value{key.removeprefix("key")}' for key, value in collected_entries) + + +async def test_iterate_entries_empty_kvs(kvs: KeyValueStore) -> None: + """Test that iterating over an empty key-value store yields nothing.""" + collected_entries = [entry async for entry in kvs.iterate_entries()] + + assert collected_entries == [] + + +async def test_iterate_entries_skips_records_deleted_during_iteration(kvs: KeyValueStore) -> None: + """Test that a record deleted after its key was listed but before its value was read is skipped.""" + for i in range(5): + await kvs.set_value(f'key{i}', f'value{i}') + + all_keys = [metadata.key for metadata in await kvs.list_keys()] + deleted_key = all_keys[-1] + + collected_entries = [] + async for key, value in kvs.iterate_entries(): + if key == all_keys[0]: + await kvs.delete_value(deleted_key) + collected_entries.append((key, value)) + + assert len(collected_entries) == 4 + assert deleted_key not in dict(collected_entries) + + +async def test_async_iteration(kvs: KeyValueStore) -> None: + """Test that the key-value store can be used directly in an `async for` loop, yielding (key, value) pairs.""" + await kvs.set_value('key1', 'value1') + await kvs.set_value('key2', 'value2') + + collected_entries = {key: value async for key, value in kvs} + + assert collected_entries == {'key1': 'value1', 'key2': 'value2'} + + async def test_drop( storage_client: StorageClient, ) -> None: From 2deace2416ce4b81ed44a7619a6af6ae5928694b Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Thu, 1 Oct 2026 12:28:27 +0000 Subject: [PATCH 02/22] feat: make KeyValueStore record iteration overridable by storage clients Add a non-abstract `KeyValueStoreClient.iterate_entries` to the storage client base class, yielding `KeyValueStoreRecord`s. The default implementation lists the keys with `iterate_keys` and reads each value with `get_value` as the iteration advances, so existing and third-party clients need no changes. Backends that can read keys together with their values more efficiently can override it. `KeyValueStore.iterate_entries` and `iterate_values` now delegate to the client method instead of combining `iterate_keys` and `get_value` themselves. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../_base/_key_value_store_client.py | 21 +++++++++ src/crawlee/storages/_key_value_store.py | 19 ++++---- tests/unit/storages/test_key_value_store.py | 45 ++++++++++++++++++- 3 files changed, 74 insertions(+), 11 deletions(-) diff --git a/src/crawlee/storage_clients/_base/_key_value_store_client.py b/src/crawlee/storage_clients/_base/_key_value_store_client.py index 2acb11c57c..c80d541750 100644 --- a/src/crawlee/storage_clients/_base/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_base/_key_value_store_client.py @@ -79,6 +79,27 @@ async def iterate_keys( if False: yield KeyValueStoreRecordMetadata() + async def iterate_entries( + self, + *, + exclusive_start_key: str | None = None, + limit: int | None = None, + ) -> AsyncIterator[KeyValueStoreRecord]: + """Iterate over all the existing records in the key-value store, including their values. + + The backend method for the `KeyValueStore.iterate_entries` and `KeyValueStore.iterate_values` calls. + + The default implementation lists the keys with `iterate_keys` and reads each value separately with + `get_value` as the iteration advances, so only a single value is held in memory at a time. A record deleted + after its key was listed but before its value was read is skipped. Backends that can read the keys together + with their values more efficiently should override this method. + """ + async for metadata in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit): + record = await self.get_value(key=metadata.key) + if record is None: + continue + yield record + @abstractmethod async def get_public_url(self, *, key: str) -> str: """Get the public URL for the given key. diff --git a/src/crawlee/storages/_key_value_store.py b/src/crawlee/storages/_key_value_store.py index aad9aef674..11d5e6983d 100644 --- a/src/crawlee/storages/_key_value_store.py +++ b/src/crawlee/storages/_key_value_store.py @@ -220,8 +220,9 @@ async def iterate_values( ) -> AsyncIterator[Any]: """Iterate over the values of the existing records in the KVS. - The values are fetched one by one as the iteration advances, so only a single value is held in memory at - a time. On remote backends this means one request per record on top of the paginated key listing. + The records are fetched lazily as the iteration advances, so only a single value is held in memory at + a time. On remote backends this means one request per record on top of the paginated key listing, unless + the storage client provides a more efficient implementation. Args: exclusive_start_key: Key to start the iteration from. @@ -240,9 +241,10 @@ async def iterate_entries( ) -> AsyncIterator[tuple[str, Any]]: """Iterate over the existing records in the KVS as `(key, value)` pairs. - The values are fetched one by one as the iteration advances, so only a single value is held in memory at - a time. On remote backends this means one request per record on top of the paginated key listing. A record - deleted after its key was listed but before its value was fetched is skipped. + The records are fetched lazily as the iteration advances, so only a single value is held in memory at + a time. On remote backends this means one request per record on top of the paginated key listing, unless + the storage client provides a more efficient implementation. A record deleted after its key was listed but + before its value was fetched is skipped. Args: exclusive_start_key: Key to start the iteration from. @@ -251,11 +253,8 @@ async def iterate_entries( Yields: A `(key, value)` tuple for each record. """ - async for metadata in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit): - record = await self._client.get_value(key=metadata.key) - if record is None: - continue - yield metadata.key, record.value + async for record in self._client.iterate_entries(exclusive_start_key=exclusive_start_key, limit=limit): + yield record.key, record.value def __aiter__(self) -> AsyncIterator[tuple[str, Any]]: """Iterate over all records in the KVS as `(key, value)` pairs. diff --git a/tests/unit/storages/test_key_value_store.py b/tests/unit/storages/test_key_value_store.py index fc2d4a8f30..e04d530871 100644 --- a/tests/unit/storages/test_key_value_store.py +++ b/tests/unit/storages/test_key_value_store.py @@ -9,11 +9,13 @@ from crawlee import service_locator from crawlee.configuration import Configuration from crawlee.storage_clients import FileSystemStorageClient, MemoryStorageClient, SqlStorageClient, StorageClient +from crawlee.storage_clients._memory import MemoryKeyValueStoreClient +from crawlee.storage_clients.models import KeyValueStoreRecord from crawlee.storages import KeyValueStore from crawlee.storages._storage_instance_manager import StorageInstanceManager if TYPE_CHECKING: - from collections.abc import AsyncGenerator + from collections.abc import AsyncGenerator, AsyncIterator from pathlib import Path @@ -333,6 +335,47 @@ async def test_iterate_entries_skips_records_deleted_during_iteration(kvs: KeyVa assert deleted_key not in dict(collected_entries) +async def test_iterate_entries_uses_storage_client_implementation() -> None: + """Test that `iterate_entries` and `iterate_values` go through the storage client's `iterate_entries`. + + Storage clients can override the default key-by-key implementation with a more efficient one, so the frontend + must delegate to the client instead of combining `iterate_keys` and `get_value` itself. + """ + + class OptimizedKeyValueStoreClient(MemoryKeyValueStoreClient): + async def iterate_entries( + self, + *, + exclusive_start_key: str | None = None, + limit: int | None = None, + ) -> AsyncIterator[KeyValueStoreRecord]: + # A single-pass implementation that never touches `get_value`. + keys = sorted(k for k in self._records if exclusive_start_key is None or k > exclusive_start_key) + for key in keys[:limit]: + record = self._records[key] + yield KeyValueStoreRecord( + key=key, + value=f'optimized-{record.value}', + content_type=record.content_type, + size=record.size, + ) + + async def get_value(self, *, key: str) -> KeyValueStoreRecord | None: + raise AssertionError(f'get_value must not be called for {key!r} when the client implements iterate_entries') + + client = await OptimizedKeyValueStoreClient.open(id=None, name=None, alias=None) + kvs = KeyValueStore(client, id=(await client.get_metadata()).id, name=None) + await kvs.set_value('key1', 'value1') + await kvs.set_value('key2', 'value2') + await kvs.set_value('key3', 'value3') + + entries = [entry async for entry in kvs.iterate_entries(exclusive_start_key='key1', limit=1)] + values = [value async for value in kvs.iterate_values()] + + assert entries == [('key2', 'optimized-value2')] + assert values == ['optimized-value1', 'optimized-value2', 'optimized-value3'] + + async def test_async_iteration(kvs: KeyValueStore) -> None: """Test that the key-value store can be used directly in an `async for` loop, yielding (key, value) pairs.""" await kvs.set_value('key1', 'value1') From 3c03eae8aeb46a5eb6196458d5cddc9e194ec626 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Thu, 1 Oct 2026 12:55:43 +0000 Subject: [PATCH 03/22] perf: read key-value store entries in bulk in the SQL and Redis clients Override `iterate_entries` where the backend can do better than one value read per key: - The SQL client selects keys, metadata and values in a single streamed query, so a store with N records costs one query instead of N+1 and still holds only one row in memory at a time. - The Redis client fetches values with one HMGET per batch instead of two round trips per record. Batches are bounded by key count and by the record sizes known from the metadata hash, so large values do not pile up in memory. A record larger than the byte limit is fetched alone. Both clients share the value decoding with `get_value` through a new `_build_record` helper. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../_base/_key_value_store_client.py | 3 +- .../_redis/_key_value_store_client.py | 71 ++++++++++++++ .../_sql/_key_value_store_client.py | 71 +++++++++++--- src/crawlee/storages/_key_value_store.py | 4 +- .../_redis/test_redis_kvs_client.py | 94 ++++++++++++++++++- .../_sql/test_sql_kvs_client.py | 39 ++++++++ tests/unit/storages/test_key_value_store.py | 13 ++- 7 files changed, 277 insertions(+), 18 deletions(-) diff --git a/src/crawlee/storage_clients/_base/_key_value_store_client.py b/src/crawlee/storage_clients/_base/_key_value_store_client.py index c80d541750..0b9f090072 100644 --- a/src/crawlee/storage_clients/_base/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_base/_key_value_store_client.py @@ -92,7 +92,8 @@ async def iterate_entries( The default implementation lists the keys with `iterate_keys` and reads each value separately with `get_value` as the iteration advances, so only a single value is held in memory at a time. A record deleted after its key was listed but before its value was read is skipped. Backends that can read the keys together - with their values more efficiently should override this method. + with their values more efficiently should override this method; such an implementation must still keep the + memory usage bounded, e.g. by streaming or by batching on the known record sizes. """ async for metadata in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit): record = await self.get_value(key=metadata.key) diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index e1190492fe..970ab50602 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -51,6 +51,15 @@ class RedisKeyValueStoreClient(KeyValueStoreClient, RedisClientMixin): _CLIENT_TYPE = 'Key-value store' """Human-readable client type for error messages.""" + _ITERATE_ENTRIES_BATCH_MAX_KEYS = 100 + """Maximum number of records fetched with a single HMGET call in `iterate_entries`.""" + + _ITERATE_ENTRIES_BATCH_MAX_BYTES = 8 * 1024 * 1024 + """Maximum total size of the records fetched with a single HMGET call in `iterate_entries`. + + A single record larger than this is still fetched, but alone in its batch. + """ + def __init__(self, storage_name: str, storage_id: str, redis: Redis) -> None: """Initialize a new instance. @@ -181,6 +190,19 @@ async def get_value(self, *, key: str) -> KeyValueStoreRecord | None: # redis-py typing issue value_bytes: bytes | None = await await_redis_response(self._redis.hget(self._items_key, key)) # ty: ignore[invalid-assignment] + return self._build_record(metadata_item, value_bytes) + + @staticmethod + def _build_record( + metadata_item: KeyValueStoreRecordMetadata, value_bytes: bytes | None + ) -> KeyValueStoreRecord | None: + """Deserialize a stored value based on its content type into a record. + + Returns None, after logging a warning, when the value is missing or cannot be decoded as the content type + claims. + """ + key = metadata_item.key + if value_bytes is None: logger.warning(f'Value for key "{key}" is missing.') return None @@ -252,6 +274,55 @@ async def iterate_keys( **MetadataUpdateParams(update_accessed_at=True), ) + @override + async def iterate_entries( + self, + *, + exclusive_start_key: str | None = None, + limit: int | None = None, + ) -> AsyncIterator[KeyValueStoreRecord]: + # Fetch the values in batches with a single HMGET per batch, instead of two round trips per record as the + # default implementation does. The batches are bounded by the record sizes known from the metadata, so a store + # with large values does not load too many of them at once. + batch: list[KeyValueStoreRecordMetadata] = [] + batch_size = 0 + + async for metadata_item in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit): + item_size = metadata_item.size or 0 + if batch and ( + len(batch) >= self._ITERATE_ENTRIES_BATCH_MAX_KEYS + or batch_size + item_size > self._ITERATE_ENTRIES_BATCH_MAX_BYTES + ): + async for record in self._fetch_records(batch): + yield record + batch, batch_size = [], 0 + + batch.append(metadata_item) + batch_size += item_size + + if batch: + async for record in self._fetch_records(batch): + yield record + + async def _fetch_records(self, batch: list[KeyValueStoreRecordMetadata]) -> AsyncIterator[KeyValueStoreRecord]: + """Fetch the values of the given records with a single HMGET call and yield the deserialized records.""" + keys = [item.key for item in batch if item.content_type != 'application/x-none'] + values: list[bytes | None] = [] + if keys: + # redis-py typing issue + values = await await_redis_response(self._redis.hmget(self._items_key, keys)) # ty: ignore[invalid-assignment] + values_by_key = dict(zip(keys, values, strict=True)) + + for metadata_item in batch: + if metadata_item.content_type == 'application/x-none': + yield KeyValueStoreRecord(value=None, **metadata_item.model_dump()) + continue + + record = self._build_record(metadata_item, values_by_key.get(metadata_item.key)) + if record is None: + continue + yield record + @override async def get_public_url(self, *, key: str) -> str: raise NotImplementedError('Public URLs are not supported for memory key-value stores.') diff --git a/src/crawlee/storage_clients/_sql/_key_value_store_client.py b/src/crawlee/storage_clients/_sql/_key_value_store_client.py index 91451fbc70..6e6e399ee7 100644 --- a/src/crawlee/storage_clients/_sql/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_sql/_key_value_store_client.py @@ -202,21 +202,33 @@ async def get_value(self, *, key: str) -> KeyValueStoreRecord | None: if not record_db: return None - # Deserialize the value based on content type - value_bytes = record_db.value + return self._build_record( + key=record_db.key, + content_type=record_db.content_type, + size=record_db.size, + value_bytes=record_db.value, + ) + + @staticmethod + def _build_record( + *, key: str, content_type: str, size: int | None, value_bytes: bytes + ) -> KeyValueStoreRecord | None: + """Deserialize a stored value based on its content type into a record. + Returns None, after logging a warning, when the stored bytes cannot be decoded as the content type claims. + """ # Handle None values - if record_db.content_type == 'application/x-none': + if content_type == 'application/x-none': value = None # Handle JSON values - elif 'application/json' in record_db.content_type: + elif 'application/json' in content_type: try: value = json.loads(value_bytes.decode('utf-8')) except (json.JSONDecodeError, UnicodeDecodeError): logger.warning(f'Failed to decode JSON value for key "{key}"') return None # Handle text values - elif record_db.content_type.startswith('text/'): + elif content_type.startswith('text/'): try: value = value_bytes.decode('utf-8') except UnicodeDecodeError: @@ -226,12 +238,7 @@ async def get_value(self, *, key: str) -> KeyValueStoreRecord | None: else: value = value_bytes - return KeyValueStoreRecord( - key=record_db.key, - value=value, - content_type=record_db.content_type, - size=record_db.size, - ) + return KeyValueStoreRecord(key=key, value=value, content_type=content_type, size=size) @retry_on_error(SQLAlchemyError) @override @@ -282,6 +289,48 @@ async def iterate_keys( await self._add_buffer_record(session) + @override + async def iterate_entries( + self, + *, + exclusive_start_key: str | None = None, + limit: int | None = None, + ) -> AsyncIterator[KeyValueStoreRecord]: + # Read the values together with the keys in a single streamed query, instead of one query per record as the + # default implementation does. Streaming keeps a single row in memory at a time. + stmt = ( + select( + self._ITEM_TABLE.key, + self._ITEM_TABLE.content_type, + self._ITEM_TABLE.size, + self._ITEM_TABLE.value, + ) + .where(self._ITEM_TABLE.key_value_store_id == self._id) + .order_by(self._ITEM_TABLE.key) + ) + + if exclusive_start_key is not None: + stmt = stmt.where(self._ITEM_TABLE.key > exclusive_start_key) + + if limit is not None: + stmt = stmt.limit(limit) + + async with self.get_session(with_simple_commit=True) as session: + result = await session.stream(stmt.execution_options(stream_results=True)) + + async for row in result: + record = self._build_record( + key=row.key, + content_type=row.content_type, + size=row.size, + value_bytes=row.value, + ) + if record is None: + continue + yield record + + await self._add_buffer_record(session) + @retry_on_error(SQLAlchemyError) @override async def record_exists(self, *, key: str) -> bool: diff --git a/src/crawlee/storages/_key_value_store.py b/src/crawlee/storages/_key_value_store.py index 11d5e6983d..dcf6e7e73b 100644 --- a/src/crawlee/storages/_key_value_store.py +++ b/src/crawlee/storages/_key_value_store.py @@ -243,8 +243,8 @@ async def iterate_entries( The records are fetched lazily as the iteration advances, so only a single value is held in memory at a time. On remote backends this means one request per record on top of the paginated key listing, unless - the storage client provides a more efficient implementation. A record deleted after its key was listed but - before its value was fetched is skipped. + the storage client provides a more efficient implementation. A record deleted while the iteration is in + progress may or may not be yielded, depending on whether its value was already read. Args: exclusive_start_key: Key to start the iteration from. diff --git a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py index b45dbbc973..d6c75e3d97 100644 --- a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py +++ b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py @@ -2,7 +2,7 @@ import asyncio import json -from typing import TYPE_CHECKING +from typing import TYPE_CHECKING, Any from unittest.mock import AsyncMock, MagicMock, patch import pytest @@ -12,7 +12,7 @@ from crawlee.storage_clients._redis._utils import await_redis_response if TYPE_CHECKING: - from collections.abc import AsyncGenerator + from collections.abc import AsyncGenerator, Iterator from fakeredis import FakeAsyncRedis @@ -257,3 +257,93 @@ async def test_set_value_does_not_retry_on_unexpected_exception(kvs_client: Redi # Verify that retry logic was not attempted assert mock_sleep.call_count == 0 + + +@pytest.fixture +def hmget_calls(kvs_client: RedisKeyValueStoreClient) -> Iterator[list[list[str]]]: + """Record the keys of every Redis `hmget` call made through the client, while still performing the call.""" + calls: list[list[str]] = [] + original_hmget = kvs_client.redis.hmget + + def recording_hmget(name: str, keys: list[str], *args: str) -> Any: + calls.append(list(keys)) + return original_hmget(name, keys, *args) + + with patch.object(kvs_client.redis, 'hmget', side_effect=recording_hmget): + yield calls + + +async def test_iterate_entries_reads_values_in_batches( + kvs_client: RedisKeyValueStoreClient, hmget_calls: list[list[str]] +) -> None: + """Test that `iterate_entries` fetches values with batched HMGET calls instead of `get_value` per key.""" + await kvs_client.set_value(key='a-json', value={'nested': [1, 2]}) + await kvs_client.set_value(key='b-text', value='plain text') + await kvs_client.set_value(key='c-bytes', value=b'\x00\x01binary', content_type='application/octet-stream') + await kvs_client.set_value(key='d-none', value=None) + + with patch.object(kvs_client, 'get_value', side_effect=AssertionError('get_value must not be called')): + records = [record async for record in kvs_client.iterate_entries()] + + assert [record.key for record in records] == ['a-json', 'b-text', 'c-bytes', 'd-none'] + assert [record.value for record in records] == [{'nested': [1, 2]}, 'plain text', b'\x00\x01binary', None] + assert records[0].content_type.startswith('application/json') + assert records[1].content_type.startswith('text/plain') + assert records[2].content_type == 'application/octet-stream' + + # All values fit in a single batch, and the `None` record needs no value fetch at all. + assert hmget_calls == [['a-json', 'b-text', 'c-bytes']] + + +async def test_iterate_entries_batches_are_bounded_by_key_count( + kvs_client: RedisKeyValueStoreClient, hmget_calls: list[list[str]] +) -> None: + """Test that `iterate_entries` splits the HMGET calls when a batch reaches the maximum number of keys.""" + for i in range(5): + await kvs_client.set_value(key=f'key{i}', value=f'value{i}') + + with patch.object(type(kvs_client), '_ITERATE_ENTRIES_BATCH_MAX_KEYS', 2): + records = [record async for record in kvs_client.iterate_entries()] + + assert [(record.key, record.value) for record in records] == [(f'key{i}', f'value{i}') for i in range(5)] + assert hmget_calls == [['key0', 'key1'], ['key2', 'key3'], ['key4']] + + +async def test_iterate_entries_batches_are_bounded_by_size( + kvs_client: RedisKeyValueStoreClient, hmget_calls: list[list[str]] +) -> None: + """Test that `iterate_entries` splits the HMGET calls by the record sizes known from the metadata. + + A record larger than the limit is still fetched, but alone in its batch. + """ + await kvs_client.set_value(key='small1', value='ab') + await kvs_client.set_value(key='small2', value='cd') + await kvs_client.set_value(key='large', value='x' * 100) + await kvs_client.set_value(key='small3', value='ef') + + with patch.object(type(kvs_client), '_ITERATE_ENTRIES_BATCH_MAX_BYTES', 10): + records = [record async for record in kvs_client.iterate_entries()] + + assert [record.key for record in records] == ['large', 'small1', 'small2', 'small3'] + assert hmget_calls == [['large'], ['small1', 'small2', 'small3']] + + +async def test_iterate_entries_with_exclusive_start_key_and_limit(kvs_client: RedisKeyValueStoreClient) -> None: + """Test that `iterate_entries` applies `exclusive_start_key` and `limit`.""" + for i in range(6): + await kvs_client.set_value(key=f'key{i}', value=f'value{i}') + + with patch.object(kvs_client, 'get_value', side_effect=AssertionError('get_value must not be called')): + records = [record async for record in kvs_client.iterate_entries(exclusive_start_key='key1', limit=3)] + + assert [(record.key, record.value) for record in records] == [ + ('key2', 'value2'), + ('key3', 'value3'), + ('key4', 'value4'), + ] + + +async def test_iterate_entries_empty_store(kvs_client: RedisKeyValueStoreClient) -> None: + records = [record async for record in kvs_client.iterate_entries()] + + assert records == [] diff --git a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py index c09d8a13dc..48fdb55b97 100644 --- a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py +++ b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py @@ -325,3 +325,42 @@ async def test_set_value_does_not_retry_on_unexpected_exception(kvs_client: SqlK # Verify that retry logic was not attempted assert mock_sleep.call_count == 0 + + +async def test_iterate_entries_reads_values_in_a_single_query(kvs_client: SqlKeyValueStoreClient) -> None: + """Test that `iterate_entries` reads keys and values together instead of calling `get_value` per key.""" + await kvs_client.set_value(key='a-json', value={'nested': [1, 2]}) + await kvs_client.set_value(key='b-text', value='plain text') + await kvs_client.set_value(key='c-bytes', value=b'\x00\x01binary', content_type='application/octet-stream') + await kvs_client.set_value(key='d-none', value=None) + + with patch.object(kvs_client, 'get_value', side_effect=AssertionError('get_value must not be called')): + records = [record async for record in kvs_client.iterate_entries()] + + assert [record.key for record in records] == ['a-json', 'b-text', 'c-bytes', 'd-none'] + assert [record.value for record in records] == [{'nested': [1, 2]}, 'plain text', b'\x00\x01binary', None] + assert records[0].content_type.startswith('application/json') + assert records[1].content_type.startswith('text/plain') + assert records[2].content_type == 'application/octet-stream' + assert all(record.size is not None for record in records) + + +async def test_iterate_entries_with_exclusive_start_key_and_limit(kvs_client: SqlKeyValueStoreClient) -> None: + """Test that `iterate_entries` applies `exclusive_start_key` and `limit` in the query.""" + for i in range(6): + await kvs_client.set_value(key=f'key{i}', value=f'value{i}') + + with patch.object(kvs_client, 'get_value', side_effect=AssertionError('get_value must not be called')): + records = [record async for record in kvs_client.iterate_entries(exclusive_start_key='key1', limit=3)] + + assert [(record.key, record.value) for record in records] == [ + ('key2', 'value2'), + ('key3', 'value3'), + ('key4', 'value4'), + ] + + +async def test_iterate_entries_empty_store(kvs_client: SqlKeyValueStoreClient) -> None: + records = [record async for record in kvs_client.iterate_entries()] + + assert records == [] diff --git a/tests/unit/storages/test_key_value_store.py b/tests/unit/storages/test_key_value_store.py index e04d530871..d8a8e2b627 100644 --- a/tests/unit/storages/test_key_value_store.py +++ b/tests/unit/storages/test_key_value_store.py @@ -317,8 +317,17 @@ async def test_iterate_entries_empty_kvs(kvs: KeyValueStore) -> None: assert collected_entries == [] -async def test_iterate_entries_skips_records_deleted_during_iteration(kvs: KeyValueStore) -> None: - """Test that a record deleted after its key was listed but before its value was read is skipped.""" +async def test_iterate_entries_skips_records_deleted_during_iteration( + kvs: KeyValueStore, storage_client: StorageClient +) -> None: + """Test that a record deleted after its key was listed but before its value was read is skipped. + + This only holds for storage clients using the default key-by-key `iterate_entries` implementation. Clients that + read the values in batches or in a single query may have already read the record when it gets deleted. + """ + if not isinstance(storage_client, (MemoryStorageClient, FileSystemStorageClient)): + pytest.skip('Storage client reads values in batches, so a deleted record may already be read.') + for i in range(5): await kvs.set_value(f'key{i}', f'value{i}') From 125bba2a7fb1f441c338f1485c11c946a60b21f6 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Fri, 2 Oct 2026 07:09:40 +0000 Subject: [PATCH 04/22] test: drop the test for records deleted during key-value store iteration Whether a record deleted mid-iteration is yielded depends on how the storage client reads values, which is not part of the contract. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- tests/unit/storages/test_key_value_store.py | 27 --------------------- 1 file changed, 27 deletions(-) diff --git a/tests/unit/storages/test_key_value_store.py b/tests/unit/storages/test_key_value_store.py index d8a8e2b627..3a8aa0d9c0 100644 --- a/tests/unit/storages/test_key_value_store.py +++ b/tests/unit/storages/test_key_value_store.py @@ -317,33 +317,6 @@ async def test_iterate_entries_empty_kvs(kvs: KeyValueStore) -> None: assert collected_entries == [] -async def test_iterate_entries_skips_records_deleted_during_iteration( - kvs: KeyValueStore, storage_client: StorageClient -) -> None: - """Test that a record deleted after its key was listed but before its value was read is skipped. - - This only holds for storage clients using the default key-by-key `iterate_entries` implementation. Clients that - read the values in batches or in a single query may have already read the record when it gets deleted. - """ - if not isinstance(storage_client, (MemoryStorageClient, FileSystemStorageClient)): - pytest.skip('Storage client reads values in batches, so a deleted record may already be read.') - - for i in range(5): - await kvs.set_value(f'key{i}', f'value{i}') - - all_keys = [metadata.key for metadata in await kvs.list_keys()] - deleted_key = all_keys[-1] - - collected_entries = [] - async for key, value in kvs.iterate_entries(): - if key == all_keys[0]: - await kvs.delete_value(deleted_key) - collected_entries.append((key, value)) - - assert len(collected_entries) == 4 - assert deleted_key not in dict(collected_entries) - - async def test_iterate_entries_uses_storage_client_implementation() -> None: """Test that `iterate_entries` and `iterate_values` go through the storage client's `iterate_entries`. From 668c170b08e2a7990688c7f1daef24b517b25004 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Fri, 2 Oct 2026 07:22:35 +0000 Subject: [PATCH 05/22] refactor: decode None values in one place in the Redis key-value store client `set_value` stores an empty byte string for `None`, so the batched read can fetch those records like any other and let `_build_record` map the `application/x-none` content type to `None`. This drops the key filtering and the per-key lookup dict from the batch fetch. Also reword the base `iterate_entries` docstring so it does not present the handling of records deleted mid-iteration as a contract. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../_base/_key_value_store_client.py | 8 +++---- .../_redis/_key_value_store_client.py | 22 ++++++++----------- .../_redis/test_redis_kvs_client.py | 4 ++-- 3 files changed, 15 insertions(+), 19 deletions(-) diff --git a/src/crawlee/storage_clients/_base/_key_value_store_client.py b/src/crawlee/storage_clients/_base/_key_value_store_client.py index 0b9f090072..80e73f7af5 100644 --- a/src/crawlee/storage_clients/_base/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_base/_key_value_store_client.py @@ -90,10 +90,10 @@ async def iterate_entries( The backend method for the `KeyValueStore.iterate_entries` and `KeyValueStore.iterate_values` calls. The default implementation lists the keys with `iterate_keys` and reads each value separately with - `get_value` as the iteration advances, so only a single value is held in memory at a time. A record deleted - after its key was listed but before its value was read is skipped. Backends that can read the keys together - with their values more efficiently should override this method; such an implementation must still keep the - memory usage bounded, e.g. by streaming or by batching on the known record sizes. + `get_value` as the iteration advances, so only a single value is held in memory at a time. Whether a record + deleted while the iteration is in progress is yielded depends on the implementation. Backends that can read + the keys together with their values more efficiently should override this method; such an implementation + must still keep the memory usage bounded, e.g. by streaming or by batching on the known record sizes. """ async for metadata in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit): record = await self.get_value(key=metadata.key) diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index 970ab50602..da57422daa 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -207,8 +207,11 @@ def _build_record( logger.warning(f'Value for key "{key}" is missing.') return None + # Handle None values + if metadata_item.content_type == 'application/x-none': + value = None # Handle JSON values - if 'application/json' in metadata_item.content_type: + elif 'application/json' in metadata_item.content_type: try: value = json.loads(value_bytes.decode('utf-8')) except (json.JSONDecodeError, UnicodeDecodeError): @@ -306,19 +309,12 @@ async def iterate_entries( async def _fetch_records(self, batch: list[KeyValueStoreRecordMetadata]) -> AsyncIterator[KeyValueStoreRecord]: """Fetch the values of the given records with a single HMGET call and yield the deserialized records.""" - keys = [item.key for item in batch if item.content_type != 'application/x-none'] - values: list[bytes | None] = [] - if keys: - # redis-py typing issue - values = await await_redis_response(self._redis.hmget(self._items_key, keys)) # ty: ignore[invalid-assignment] - values_by_key = dict(zip(keys, values, strict=True)) - - for metadata_item in batch: - if metadata_item.content_type == 'application/x-none': - yield KeyValueStoreRecord(value=None, **metadata_item.model_dump()) - continue + keys = [item.key for item in batch] + # redis-py typing issue + values: list[bytes | None] = await await_redis_response(self._redis.hmget(self._items_key, keys)) # ty: ignore[invalid-assignment] - record = self._build_record(metadata_item, values_by_key.get(metadata_item.key)) + for metadata_item, value_bytes in zip(batch, values, strict=True): + record = self._build_record(metadata_item, value_bytes) if record is None: continue yield record diff --git a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py index d6c75e3d97..d6102681dd 100644 --- a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py +++ b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py @@ -291,8 +291,8 @@ async def test_iterate_entries_reads_values_in_batches( assert records[1].content_type.startswith('text/plain') assert records[2].content_type == 'application/octet-stream' - # All values fit in a single batch, and the `None` record needs no value fetch at all. - assert hmget_calls == [['a-json', 'b-text', 'c-bytes']] + # All values fit in a single batch. + assert hmget_calls == [['a-json', 'b-text', 'c-bytes', 'd-none']] async def test_iterate_entries_batches_are_bounded_by_key_count( From 0e255a5de1ded433a5e076a0f49326f105f537c3 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Fri, 2 Oct 2026 07:29:58 +0000 Subject: [PATCH 06/22] docs: describe the bulk iterate_entries overrides in docstrings Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../storage_clients/_redis/_key_value_store_client.py | 9 ++++++--- .../storage_clients/_sql/_key_value_store_client.py | 7 +++++-- 2 files changed, 11 insertions(+), 5 deletions(-) diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index da57422daa..b0129b3d0e 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -284,9 +284,12 @@ async def iterate_entries( exclusive_start_key: str | None = None, limit: int | None = None, ) -> AsyncIterator[KeyValueStoreRecord]: - # Fetch the values in batches with a single HMGET per batch, instead of two round trips per record as the - # default implementation does. The batches are bounded by the record sizes known from the metadata, so a store - # with large values does not load too many of them at once. + """Iterate over all the existing records in the key-value store, including their values. + + The values are fetched in batches with a single HMGET call per batch, instead of two round trips per record + as the default implementation does. The batches are bounded by the record sizes known from the metadata, so + a store with large values does not load too many of them at once. + """ batch: list[KeyValueStoreRecordMetadata] = [] batch_size = 0 diff --git a/src/crawlee/storage_clients/_sql/_key_value_store_client.py b/src/crawlee/storage_clients/_sql/_key_value_store_client.py index 6e6e399ee7..43d166dd9b 100644 --- a/src/crawlee/storage_clients/_sql/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_sql/_key_value_store_client.py @@ -296,8 +296,11 @@ async def iterate_entries( exclusive_start_key: str | None = None, limit: int | None = None, ) -> AsyncIterator[KeyValueStoreRecord]: - # Read the values together with the keys in a single streamed query, instead of one query per record as the - # default implementation does. Streaming keeps a single row in memory at a time. + """Iterate over all the existing records in the key-value store, including their values. + + The values are read together with the keys in a single streamed query, instead of one query per record as + the default implementation does. Streaming keeps a single row in memory at a time. + """ stmt = ( select( self._ITEM_TABLE.key, From 6c0cc3fba89f5d57c8c3f5353cb37b08509dd567 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Fri, 2 Oct 2026 07:50:56 +0000 Subject: [PATCH 07/22] refactor: reject a decoding Redis client instead of suppressing the type checker redis-py types every reply as `bytes | str | None` because a client created with `decode_responses=True` returns strings. The key-value store client stores binary values and needs the raw bytes back, so it narrowed the type with `ty: ignore` and would fail with an `AttributeError` on such a client. Add `expect_bytes`, which narrows a reply to bytes and raises a `TypeError` naming the unsupported option, use it at both value-reading sites, and document the requirement on `RedisStorageClient`. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../_redis/_key_value_store_client.py | 8 +++---- .../storage_clients/_redis/_storage_client.py | 3 ++- src/crawlee/storage_clients/_redis/_utils.py | 24 +++++++++++++++++++ .../_redis/test_redis_kvs_client.py | 16 +++++++++++-- 4 files changed, 43 insertions(+), 8 deletions(-) diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index b0129b3d0e..3f2c5c2231 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -13,7 +13,7 @@ from crawlee.storage_clients.models import KeyValueStoreMetadata, KeyValueStoreRecord, KeyValueStoreRecordMetadata from ._client_mixin import MetadataUpdateParams, RedisClientMixin -from ._utils import await_redis_response +from ._utils import await_redis_response, expect_bytes if TYPE_CHECKING: from collections.abc import AsyncIterator @@ -187,8 +187,7 @@ async def get_value(self, *, key: str) -> KeyValueStoreRecord | None: return KeyValueStoreRecord(value=None, **metadata_item.model_dump()) # Query the record by key - # redis-py typing issue - value_bytes: bytes | None = await await_redis_response(self._redis.hget(self._items_key, key)) # ty: ignore[invalid-assignment] + value_bytes = expect_bytes(await await_redis_response(self._redis.hget(self._items_key, key))) return self._build_record(metadata_item, value_bytes) @@ -313,8 +312,7 @@ async def iterate_entries( async def _fetch_records(self, batch: list[KeyValueStoreRecordMetadata]) -> AsyncIterator[KeyValueStoreRecord]: """Fetch the values of the given records with a single HMGET call and yield the deserialized records.""" keys = [item.key for item in batch] - # redis-py typing issue - values: list[bytes | None] = await await_redis_response(self._redis.hmget(self._items_key, keys)) # ty: ignore[invalid-assignment] + values = expect_bytes(await await_redis_response(self._redis.hmget(self._items_key, keys))) for metadata_item, value_bytes in zip(batch, values, strict=True): record = self._build_record(metadata_item, value_bytes) diff --git a/src/crawlee/storage_clients/_redis/_storage_client.py b/src/crawlee/storage_clients/_redis/_storage_client.py index a6c39f5def..9121c28db9 100644 --- a/src/crawlee/storage_clients/_redis/_storage_client.py +++ b/src/crawlee/storage_clients/_redis/_storage_client.py @@ -23,7 +23,8 @@ class RedisStorageClient(StorageClient): to a Redis database v8.0+. Each storage type uses Redis-specific data structures and key patterns for efficient storage and retrieval. - The client accepts either a Redis connection string or a pre-configured Redis client instance. + The client accepts either a Redis connection string or a pre-configured Redis client instance. The Redis client + must return raw bytes, which is the default; a client created with `decode_responses=True` is not supported. Exactly one of these parameters must be provided during initialization. Storage types use the following Redis data structures: diff --git a/src/crawlee/storage_clients/_redis/_utils.py b/src/crawlee/storage_clients/_redis/_utils.py index 92bd1afce1..a6fb9e331f 100644 --- a/src/crawlee/storage_clients/_redis/_utils.py +++ b/src/crawlee/storage_clients/_redis/_utils.py @@ -18,6 +18,30 @@ async def await_redis_response(response: Awaitable[T] | T) -> T: return response +@overload +def expect_bytes(value: bytes | str | None) -> bytes | None: ... +@overload +def expect_bytes(value: list[bytes | str | None]) -> list[bytes | None]: ... + + +def expect_bytes(value: bytes | str | list[bytes | str | None] | None) -> bytes | list[bytes | None] | None: + """Narrow a Redis reply to raw bytes, rejecting a client that decodes responses. + + redis-py types every reply as `bytes | str | None`, because a client created with `decode_responses=True` returns + strings. The storage clients store binary values and need the raw bytes back, so such a client is not supported. + + Raises: + TypeError: If the reply contains a string, i.e. the Redis client decodes responses. + """ + values = value if isinstance(value, list) else [value] + if any(isinstance(item, str) for item in values): + raise TypeError( + 'The Redis client returned a decoded string instead of raw bytes. The Redis storage client requires ' + 'a Redis client created without `decode_responses=True`.' + ) + return cast('bytes | list[bytes | None] | None', value) + + def read_lua_script(script_name: str) -> str: """Read a Lua script from a file.""" file_path = Path(__file__).parent / 'lua_scripts' / script_name diff --git a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py index d6102681dd..3391fc65a9 100644 --- a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py +++ b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py @@ -6,6 +6,7 @@ from unittest.mock import AsyncMock, MagicMock, patch import pytest +from fakeredis import FakeAsyncRedis from redis.exceptions import RedisError from crawlee.storage_clients import RedisStorageClient @@ -14,8 +15,6 @@ if TYPE_CHECKING: from collections.abc import AsyncGenerator, Iterator - from fakeredis import FakeAsyncRedis - from crawlee.storage_clients._redis import RedisKeyValueStoreClient @@ -347,3 +346,16 @@ async def test_iterate_entries_empty_store(kvs_client: RedisKeyValueStoreClient) records = [record async for record in kvs_client.iterate_entries()] assert records == [] + + +async def test_decoding_redis_client_is_rejected(suppress_user_warning: None) -> None: # noqa: ARG001 + """Test that reading values through a Redis client created with `decode_responses=True` raises a clear error.""" + storage_client = RedisStorageClient(redis=FakeAsyncRedis(decode_responses=True)) + kvs_client = await storage_client.create_kvs_client(name='decoding_kvs') + await kvs_client.set_value(key='key', value='value') + + with pytest.raises(TypeError, match='decode_responses'): + await kvs_client.get_value(key='key') + + with pytest.raises(TypeError, match='decode_responses'): + _ = [record async for record in kvs_client.iterate_entries()] From 95a5154170e2e7eadb2a1879162e5049b47bb562 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Fri, 2 Oct 2026 07:57:52 +0000 Subject: [PATCH 08/22] refactor: narrow Redis replies without a cast `expect_bytes` now handles a single reply, which `isinstance` narrows on its own, and the HMGET result is narrowed with a comprehension. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../_redis/_key_value_store_client.py | 2 +- src/crawlee/storage_clients/_redis/_utils.py | 15 ++++----------- 2 files changed, 5 insertions(+), 12 deletions(-) diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index 3f2c5c2231..7e54e1e523 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -312,7 +312,7 @@ async def iterate_entries( async def _fetch_records(self, batch: list[KeyValueStoreRecordMetadata]) -> AsyncIterator[KeyValueStoreRecord]: """Fetch the values of the given records with a single HMGET call and yield the deserialized records.""" keys = [item.key for item in batch] - values = expect_bytes(await await_redis_response(self._redis.hmget(self._items_key, keys))) + values = [expect_bytes(v) for v in await await_redis_response(self._redis.hmget(self._items_key, keys))] for metadata_item, value_bytes in zip(batch, values, strict=True): record = self._build_record(metadata_item, value_bytes) diff --git a/src/crawlee/storage_clients/_redis/_utils.py b/src/crawlee/storage_clients/_redis/_utils.py index a6fb9e331f..8d690ace7b 100644 --- a/src/crawlee/storage_clients/_redis/_utils.py +++ b/src/crawlee/storage_clients/_redis/_utils.py @@ -18,28 +18,21 @@ async def await_redis_response(response: Awaitable[T] | T) -> T: return response -@overload -def expect_bytes(value: bytes | str | None) -> bytes | None: ... -@overload -def expect_bytes(value: list[bytes | str | None]) -> list[bytes | None]: ... - - -def expect_bytes(value: bytes | str | list[bytes | str | None] | None) -> bytes | list[bytes | None] | None: +def expect_bytes(value: bytes | str | None) -> bytes | None: """Narrow a Redis reply to raw bytes, rejecting a client that decodes responses. redis-py types every reply as `bytes | str | None`, because a client created with `decode_responses=True` returns strings. The storage clients store binary values and need the raw bytes back, so such a client is not supported. Raises: - TypeError: If the reply contains a string, i.e. the Redis client decodes responses. + TypeError: If the reply is a string, i.e. the Redis client decodes responses. """ - values = value if isinstance(value, list) else [value] - if any(isinstance(item, str) for item in values): + if isinstance(value, str): raise TypeError( 'The Redis client returned a decoded string instead of raw bytes. The Redis storage client requires ' 'a Redis client created without `decode_responses=True`.' ) - return cast('bytes | list[bytes | None] | None', value) + return value def read_lua_script(script_name: str) -> str: From a07f81aa92306cf3e72653f6b2aea24d881af853 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Fri, 2 Oct 2026 08:08:15 +0000 Subject: [PATCH 09/22] docs: scope the raw-bytes requirement to the Redis key-value store client Datasets use RedisJSON and request queues parse JSON strings, so both work with a decoding Redis client. Only key-value store values are binary. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../storage_clients/_redis/_key_value_store_client.py | 3 +++ src/crawlee/storage_clients/_redis/_storage_client.py | 5 +++-- src/crawlee/storage_clients/_redis/_utils.py | 7 ++++--- 3 files changed, 10 insertions(+), 5 deletions(-) diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index 7e54e1e523..d98f04229f 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -40,6 +40,9 @@ class RedisKeyValueStoreClient(KeyValueStoreClient, RedisClientMixin): All operations are atomic through Redis hash operations and pipeline transactions. The client supports concurrent access through Redis's built-in atomic operations for hash fields. + + Values are stored as raw bytes, so the Redis client must not decode responses. A Redis client created with + `decode_responses=True` is rejected when a value is read. """ _DEFAULT_NAME = 'default' diff --git a/src/crawlee/storage_clients/_redis/_storage_client.py b/src/crawlee/storage_clients/_redis/_storage_client.py index 9121c28db9..3450df1390 100644 --- a/src/crawlee/storage_clients/_redis/_storage_client.py +++ b/src/crawlee/storage_clients/_redis/_storage_client.py @@ -23,8 +23,9 @@ class RedisStorageClient(StorageClient): to a Redis database v8.0+. Each storage type uses Redis-specific data structures and key patterns for efficient storage and retrieval. - The client accepts either a Redis connection string or a pre-configured Redis client instance. The Redis client - must return raw bytes, which is the default; a client created with `decode_responses=True` is not supported. + The client accepts either a Redis connection string or a pre-configured Redis client instance. Key-value stores + hold binary values, so their client needs a Redis client that returns raw bytes, which is the default; a client + created with `decode_responses=True` works for datasets and request queues only. Exactly one of these parameters must be provided during initialization. Storage types use the following Redis data structures: diff --git a/src/crawlee/storage_clients/_redis/_utils.py b/src/crawlee/storage_clients/_redis/_utils.py index 8d690ace7b..5359191d74 100644 --- a/src/crawlee/storage_clients/_redis/_utils.py +++ b/src/crawlee/storage_clients/_redis/_utils.py @@ -22,15 +22,16 @@ def expect_bytes(value: bytes | str | None) -> bytes | None: """Narrow a Redis reply to raw bytes, rejecting a client that decodes responses. redis-py types every reply as `bytes | str | None`, because a client created with `decode_responses=True` returns - strings. The storage clients store binary values and need the raw bytes back, so such a client is not supported. + strings. The key-value store client stores binary values and needs the raw bytes back, so it rejects such + a client. Raises: TypeError: If the reply is a string, i.e. the Redis client decodes responses. """ if isinstance(value, str): raise TypeError( - 'The Redis client returned a decoded string instead of raw bytes. The Redis storage client requires ' - 'a Redis client created without `decode_responses=True`.' + 'The Redis client returned a decoded string instead of raw bytes. The Redis key-value store client ' + 'requires a Redis client created without `decode_responses=True`.' ) return value From 1f32875d2937e0ce4e7d4a531febff336857fd38 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Fri, 2 Oct 2026 08:33:45 +0000 Subject: [PATCH 10/22] fix: bound memory of the SQL key-value store iteration with keyset pages `stream_results` buffers up to 1000 rows client side, so iterating a store of large values held hundreds of megabytes before the first yield, and the single open transaction lasted for the whole iteration. The SQL client now reads record metadata in keyset-paginated pages, splits each page into batches bounded by the record sizes, and fetches every batch's values with one query. Each query runs in its own short session. The batching rule is shared with the Redis client through a new `batch_records_by_size` helper. Measured with 300 records of 1 MiB: peak memory drops from 300 MiB to 9 MiB. With 2000 small records the iteration issues 41 selects instead of 2001 and is about 40 times faster than the default implementation. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../_redis/_key_value_store_client.py | 27 ++--- .../_sql/_key_value_store_client.py | 99 ++++++++++++++----- src/crawlee/storage_clients/_utils.py | 36 +++++++ .../_sql/test_sql_kvs_client.py | 55 ++++++++++- 4 files changed, 174 insertions(+), 43 deletions(-) create mode 100644 src/crawlee/storage_clients/_utils.py diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index d98f04229f..e1afb76ff1 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -10,6 +10,7 @@ from crawlee._utils.file import infer_mime_type from crawlee._utils.retry import retry_on_error from crawlee.storage_clients._base import KeyValueStoreClient +from crawlee.storage_clients._utils import batch_records_by_size from crawlee.storage_clients.models import KeyValueStoreMetadata, KeyValueStoreRecord, KeyValueStoreRecordMetadata from ._client_mixin import MetadataUpdateParams, RedisClientMixin @@ -292,23 +293,15 @@ async def iterate_entries( as the default implementation does. The batches are bounded by the record sizes known from the metadata, so a store with large values does not load too many of them at once. """ - batch: list[KeyValueStoreRecordMetadata] = [] - batch_size = 0 - - async for metadata_item in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit): - item_size = metadata_item.size or 0 - if batch and ( - len(batch) >= self._ITERATE_ENTRIES_BATCH_MAX_KEYS - or batch_size + item_size > self._ITERATE_ENTRIES_BATCH_MAX_BYTES - ): - async for record in self._fetch_records(batch): - yield record - batch, batch_size = [], 0 - - batch.append(metadata_item) - batch_size += item_size - - if batch: + metadata_items = [ + item async for item in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit) + ] + + for batch in batch_records_by_size( + metadata_items, + max_records=self._ITERATE_ENTRIES_BATCH_MAX_KEYS, + max_bytes=self._ITERATE_ENTRIES_BATCH_MAX_BYTES, + ): async for record in self._fetch_records(batch): yield record diff --git a/src/crawlee/storage_clients/_sql/_key_value_store_client.py b/src/crawlee/storage_clients/_sql/_key_value_store_client.py index 43d166dd9b..d2d97d1dc6 100644 --- a/src/crawlee/storage_clients/_sql/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_sql/_key_value_store_client.py @@ -13,6 +13,7 @@ from crawlee._utils.file import infer_mime_type from crawlee._utils.retry import retry_on_error from crawlee.storage_clients._base import KeyValueStoreClient +from crawlee.storage_clients._utils import batch_records_by_size from crawlee.storage_clients.models import ( KeyValueStoreMetadata, KeyValueStoreRecord, @@ -59,6 +60,15 @@ class SqlKeyValueStoreClient(KeyValueStoreClient, SqlClientMixin): _DEFAULT_NAME = 'default' """Default dataset name used when no name is provided.""" + _ITERATE_ENTRIES_BATCH_MAX_KEYS = 100 + """Maximum number of records read with a single query in `iterate_entries`.""" + + _ITERATE_ENTRIES_BATCH_MAX_BYTES = 8 * 1024 * 1024 + """Maximum total size of the records read with a single query in `iterate_entries`. + + A single record larger than this is still read, but alone in its batch. + """ + _METADATA_TABLE = KeyValueStoreMetadataDb """SQLAlchemy model for key-value store metadata.""" @@ -298,9 +308,63 @@ async def iterate_entries( ) -> AsyncIterator[KeyValueStoreRecord]: """Iterate over all the existing records in the key-value store, including their values. - The values are read together with the keys in a single streamed query, instead of one query per record as - the default implementation does. Streaming keeps a single row in memory at a time. + The records are read in keyset-paginated pages of metadata, and the values of each page are then fetched with + a single query per batch, instead of one query per record as the default implementation does. The batches are + bounded by the record sizes, so a store with large values does not load too many of them at once. Every query + runs in its own short session, so no transaction stays open while the consumer processes the records. """ + last_key = exclusive_start_key + remaining = limit + + while remaining is None or remaining > 0: + page_size = self._ITERATE_ENTRIES_BATCH_MAX_KEYS + if remaining is not None: + page_size = min(page_size, remaining) + + page = await self._list_record_metadata(exclusive_start_key=last_key, limit=page_size) + if not page: + return + + for batch in batch_records_by_size( + page, + max_records=self._ITERATE_ENTRIES_BATCH_MAX_KEYS, + max_bytes=self._ITERATE_ENTRIES_BATCH_MAX_BYTES, + ): + for record in await self._fetch_records(batch): + yield record + + last_key = page[-1].key + if remaining is not None: + remaining -= len(page) + if len(page) < page_size: + return + + @retry_on_error(SQLAlchemyError) + async def _list_record_metadata( + self, *, exclusive_start_key: str | None, limit: int + ) -> list[KeyValueStoreRecordMetadata]: + """Read one page of record metadata, ordered by key and starting after `exclusive_start_key`.""" + stmt = ( + select(self._ITEM_TABLE.key, self._ITEM_TABLE.content_type, self._ITEM_TABLE.size) + .where(self._ITEM_TABLE.key_value_store_id == self._id) + .order_by(self._ITEM_TABLE.key) + .limit(limit) + ) + if exclusive_start_key is not None: + stmt = stmt.where(self._ITEM_TABLE.key > exclusive_start_key) + + async with self.get_session(with_simple_commit=True) as session: + result = await session.execute(stmt) + page = [ + KeyValueStoreRecordMetadata(key=row.key, content_type=row.content_type, size=row.size) for row in result + ] + await self._add_buffer_record(session) + + return page + + @retry_on_error(SQLAlchemyError) + async def _fetch_records(self, batch: list[KeyValueStoreRecordMetadata]) -> list[KeyValueStoreRecord]: + """Fetch the values of the given records with a single query and return the deserialized records.""" stmt = ( select( self._ITEM_TABLE.key, @@ -308,31 +372,22 @@ async def iterate_entries( self._ITEM_TABLE.size, self._ITEM_TABLE.value, ) - .where(self._ITEM_TABLE.key_value_store_id == self._id) + .where( + self._ITEM_TABLE.key_value_store_id == self._id, + self._ITEM_TABLE.key.in_([item.key for item in batch]), + ) .order_by(self._ITEM_TABLE.key) ) - if exclusive_start_key is not None: - stmt = stmt.where(self._ITEM_TABLE.key > exclusive_start_key) - - if limit is not None: - stmt = stmt.limit(limit) - async with self.get_session(with_simple_commit=True) as session: - result = await session.stream(stmt.execution_options(stream_results=True)) - - async for row in result: - record = self._build_record( - key=row.key, - content_type=row.content_type, - size=row.size, - value_bytes=row.value, - ) - if record is None: - continue - yield record + result = await session.execute(stmt) + rows = result.all() - await self._add_buffer_record(session) + records = ( + self._build_record(key=row.key, content_type=row.content_type, size=row.size, value_bytes=row.value) + for row in rows + ) + return [record for record in records if record is not None] @retry_on_error(SQLAlchemyError) @override diff --git a/src/crawlee/storage_clients/_utils.py b/src/crawlee/storage_clients/_utils.py new file mode 100644 index 0000000000..f02e285fd7 --- /dev/null +++ b/src/crawlee/storage_clients/_utils.py @@ -0,0 +1,36 @@ +from __future__ import annotations + +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from collections.abc import Iterable, Iterator + + from crawlee.storage_clients.models import KeyValueStoreRecordMetadata + + +def batch_records_by_size( + records: Iterable[KeyValueStoreRecordMetadata], + *, + max_records: int, + max_bytes: int, +) -> Iterator[list[KeyValueStoreRecordMetadata]]: + """Group record metadata into batches bounded by record count and by total record size. + + Used by storage clients that read the values of several records at once, so that a store with large values does + not load too many of them into memory in a single read. A record whose size alone exceeds `max_bytes` is still + yielded, but alone in its batch. A record with unknown size counts as empty. + """ + batch: list[KeyValueStoreRecordMetadata] = [] + batch_size = 0 + + for record in records: + record_size = record.size or 0 + if batch and (len(batch) >= max_records or batch_size + record_size > max_bytes): + yield batch + batch, batch_size = [], 0 + + batch.append(record) + batch_size += record_size + + if batch: + yield batch diff --git a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py index 48fdb55b97..5aea03570a 100644 --- a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py +++ b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py @@ -13,10 +13,10 @@ from crawlee.configuration import Configuration from crawlee.storage_clients import SqlStorageClient from crawlee.storage_clients._sql._db_models import KeyValueStoreMetadataDb, KeyValueStoreRecordDb -from crawlee.storage_clients.models import KeyValueStoreMetadata +from crawlee.storage_clients.models import KeyValueStoreMetadata, KeyValueStoreRecord, KeyValueStoreRecordMetadata if TYPE_CHECKING: - from collections.abc import AsyncGenerator + from collections.abc import AsyncGenerator, Iterator from pathlib import Path from sqlalchemy import Connection @@ -327,8 +327,22 @@ async def test_set_value_does_not_retry_on_unexpected_exception(kvs_client: SqlK assert mock_sleep.call_count == 0 -async def test_iterate_entries_reads_values_in_a_single_query(kvs_client: SqlKeyValueStoreClient) -> None: - """Test that `iterate_entries` reads keys and values together instead of calling `get_value` per key.""" +@pytest.fixture +def fetched_batches(kvs_client: SqlKeyValueStoreClient) -> Iterator[list[list[str]]]: + """Record the keys of every value batch `iterate_entries` fetches, while still performing the fetch.""" + calls: list[list[str]] = [] + original_fetch = kvs_client._fetch_records + + async def recording_fetch(batch: list[KeyValueStoreRecordMetadata]) -> list[KeyValueStoreRecord]: + calls.append([item.key for item in batch]) + return await original_fetch(batch) + + with patch.object(kvs_client, '_fetch_records', side_effect=recording_fetch): + yield calls + + +async def test_iterate_entries_reads_values_in_batches(kvs_client: SqlKeyValueStoreClient) -> None: + """Test that `iterate_entries` reads values in batches instead of calling `get_value` per key.""" await kvs_client.set_value(key='a-json', value={'nested': [1, 2]}) await kvs_client.set_value(key='b-text', value='plain text') await kvs_client.set_value(key='c-bytes', value=b'\x00\x01binary', content_type='application/octet-stream') @@ -364,3 +378,36 @@ async def test_iterate_entries_empty_store(kvs_client: SqlKeyValueStoreClient) - records = [record async for record in kvs_client.iterate_entries()] assert records == [] + + +async def test_iterate_entries_batches_are_bounded_by_key_count( + kvs_client: SqlKeyValueStoreClient, fetched_batches: list[list[str]] +) -> None: + """Test that `iterate_entries` splits the value queries when a batch reaches the maximum number of keys.""" + for i in range(5): + await kvs_client.set_value(key=f'key{i}', value=f'value{i}') + + with patch.object(type(kvs_client), '_ITERATE_ENTRIES_BATCH_MAX_KEYS', 2): + records = [record async for record in kvs_client.iterate_entries()] + + assert [(record.key, record.value) for record in records] == [(f'key{i}', f'value{i}') for i in range(5)] + assert fetched_batches == [['key0', 'key1'], ['key2', 'key3'], ['key4']] + + +async def test_iterate_entries_batches_are_bounded_by_size( + kvs_client: SqlKeyValueStoreClient, fetched_batches: list[list[str]] +) -> None: + """Test that `iterate_entries` splits the value queries by the record sizes. + + A record larger than the limit is still read, but alone in its batch. + """ + await kvs_client.set_value(key='small1', value='ab') + await kvs_client.set_value(key='small2', value='cd') + await kvs_client.set_value(key='large', value='x' * 100) + await kvs_client.set_value(key='small3', value='ef') + + with patch.object(type(kvs_client), '_ITERATE_ENTRIES_BATCH_MAX_BYTES', 10): + records = [record async for record in kvs_client.iterate_entries()] + + assert [record.key for record in records] == ['large', 'small1', 'small2', 'small3'] + assert fetched_batches == [['large'], ['small1', 'small2', 'small3']] From f809fe76f3e44bc05bd82090cb571e8836674f0a Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Fri, 2 Oct 2026 11:26:06 +0000 Subject: [PATCH 11/22] fix: retry the batched Redis value read and test the batching helper The default iteration retries every read through `get_value`, but the batched HMGET in the Redis override was not retried. It now returns a list under the same `retry_on_error` as the other Redis reads, matching the SQL client. Add a direct unit test for `batch_records_by_size`, covering the count bound, the size bound, an oversized record alone in its batch and records of unknown size. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../_redis/_key_value_store_client.py | 17 +++---- tests/unit/storage_clients/test_utils.py | 48 +++++++++++++++++++ 2 files changed, 57 insertions(+), 8 deletions(-) create mode 100644 tests/unit/storage_clients/test_utils.py diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index e1afb76ff1..a501eff587 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -302,19 +302,20 @@ async def iterate_entries( max_records=self._ITERATE_ENTRIES_BATCH_MAX_KEYS, max_bytes=self._ITERATE_ENTRIES_BATCH_MAX_BYTES, ): - async for record in self._fetch_records(batch): + for record in await self._fetch_records(batch): yield record - async def _fetch_records(self, batch: list[KeyValueStoreRecordMetadata]) -> AsyncIterator[KeyValueStoreRecord]: - """Fetch the values of the given records with a single HMGET call and yield the deserialized records.""" + @retry_on_error(RedisError) + async def _fetch_records(self, batch: list[KeyValueStoreRecordMetadata]) -> list[KeyValueStoreRecord]: + """Fetch the values of the given records with a single HMGET call and return the deserialized records.""" keys = [item.key for item in batch] values = [expect_bytes(v) for v in await await_redis_response(self._redis.hmget(self._items_key, keys))] - for metadata_item, value_bytes in zip(batch, values, strict=True): - record = self._build_record(metadata_item, value_bytes) - if record is None: - continue - yield record + records = ( + self._build_record(metadata_item, value_bytes) + for metadata_item, value_bytes in zip(batch, values, strict=True) + ) + return [record for record in records if record is not None] @override async def get_public_url(self, *, key: str) -> str: diff --git a/tests/unit/storage_clients/test_utils.py b/tests/unit/storage_clients/test_utils.py new file mode 100644 index 0000000000..2d546e8363 --- /dev/null +++ b/tests/unit/storage_clients/test_utils.py @@ -0,0 +1,48 @@ +from __future__ import annotations + +from crawlee.storage_clients._utils import batch_records_by_size +from crawlee.storage_clients.models import KeyValueStoreRecordMetadata + + +def _record(key: str, size: int | None) -> KeyValueStoreRecordMetadata: + return KeyValueStoreRecordMetadata(key=key, content_type='application/octet-stream', size=size) + + +def _keys(batches: list[list[KeyValueStoreRecordMetadata]]) -> list[list[str]]: + return [[record.key for record in batch] for batch in batches] + + +def test_batches_are_bounded_by_record_count() -> None: + records = [_record(f'k{i}', 1) for i in range(5)] + + batches = list(batch_records_by_size(records, max_records=2, max_bytes=1000)) + + assert _keys(batches) == [['k0', 'k1'], ['k2', 'k3'], ['k4']] + + +def test_batches_are_bounded_by_total_size() -> None: + records = [_record('a', 4), _record('b', 4), _record('c', 4), _record('d', 4)] + + batches = list(batch_records_by_size(records, max_records=100, max_bytes=10)) + + assert _keys(batches) == [['a', 'b'], ['c', 'd']] + + +def test_oversized_record_is_alone_in_its_batch() -> None: + records = [_record('small1', 2), _record('large', 100), _record('small2', 2), _record('small3', 2)] + + batches = list(batch_records_by_size(records, max_records=100, max_bytes=10)) + + assert _keys(batches) == [['small1'], ['large'], ['small2', 'small3']] + + +def test_unknown_size_counts_as_empty() -> None: + records = [_record('a', None), _record('b', None), _record('c', 10)] + + batches = list(batch_records_by_size(records, max_records=100, max_bytes=10)) + + assert _keys(batches) == [['a', 'b', 'c']] + + +def test_no_records_yield_no_batches() -> None: + assert list(batch_records_by_size([], max_records=100, max_bytes=10)) == [] From c0c334e10be475bcb919d21d14d950e4064c25b2 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Fri, 2 Oct 2026 11:31:16 +0000 Subject: [PATCH 12/22] test: give the batching helper test module a unique basename The test directories have no `__init__.py`, so pytest cannot import two modules named `test_utils.py` and aborted the whole unit test collection. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- .../{test_utils.py => test_batch_records_by_size.py} | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename tests/unit/storage_clients/{test_utils.py => test_batch_records_by_size.py} (100%) diff --git a/tests/unit/storage_clients/test_utils.py b/tests/unit/storage_clients/test_batch_records_by_size.py similarity index 100% rename from tests/unit/storage_clients/test_utils.py rename to tests/unit/storage_clients/test_batch_records_by_size.py From f7db136561a0cefdf09ee9490f7d4ff6d977eb4f Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 6 Oct 2026 10:36:57 +0200 Subject: [PATCH 13/22] fix: skip Redis records deleted mid-iteration without a missing-value warning --- .../_redis/_key_value_store_client.py | 2 ++ .../_redis/test_redis_kvs_client.py | 15 +++++++++++++++ 2 files changed, 17 insertions(+) diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index a501eff587..3a3ff08f96 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -311,9 +311,11 @@ async def _fetch_records(self, batch: list[KeyValueStoreRecordMetadata]) -> list keys = [item.key for item in batch] values = [expect_bytes(v) for v in await await_redis_response(self._redis.hmget(self._items_key, keys))] + # A missing value means the record was deleted after its metadata was listed, so it is skipped silently. records = ( self._build_record(metadata_item, value_bytes) for metadata_item, value_bytes in zip(batch, values, strict=True) + if value_bytes is not None ) return [record for record in records if record is not None] diff --git a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py index 3391fc65a9..af2039633c 100644 --- a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py +++ b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py @@ -342,6 +342,21 @@ async def test_iterate_entries_with_exclusive_start_key_and_limit(kvs_client: Re ] +async def test_iterate_entries_skips_value_deleted_after_listing( + kvs_client: RedisKeyValueStoreClient, caplog: pytest.LogCaptureFixture +) -> None: + """Test that `iterate_entries` silently skips a record whose value is gone by the time its batch is fetched.""" + await kvs_client.set_value(key='kept', value='a') + await kvs_client.set_value(key='removed', value='b') + await await_redis_response(kvs_client.redis.hdel(kvs_client._items_key, 'removed')) + + with caplog.at_level('WARNING'): + records = [record async for record in kvs_client.iterate_entries()] + + assert [(record.key, record.value) for record in records] == [('kept', 'a')] + assert 'missing' not in caplog.text + + async def test_iterate_entries_empty_store(kvs_client: RedisKeyValueStoreClient) -> None: records = [record async for record in kvs_client.iterate_entries()] From 5918899f07d255ebdbf80a75a1545cf2e97bf762 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 6 Oct 2026 10:36:59 +0200 Subject: [PATCH 14/22] docs: correct memory and round-trip claims in KVS iteration docstrings --- .../_redis/_key_value_store_client.py | 2 +- src/crawlee/storage_clients/_utils.py | 6 +++--- src/crawlee/storages/_key_value_store.py | 12 ++++++------ 3 files changed, 10 insertions(+), 10 deletions(-) diff --git a/src/crawlee/storage_clients/_redis/_key_value_store_client.py b/src/crawlee/storage_clients/_redis/_key_value_store_client.py index 3a3ff08f96..bad1b76234 100644 --- a/src/crawlee/storage_clients/_redis/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_redis/_key_value_store_client.py @@ -289,7 +289,7 @@ async def iterate_entries( ) -> AsyncIterator[KeyValueStoreRecord]: """Iterate over all the existing records in the key-value store, including their values. - The values are fetched in batches with a single HMGET call per batch, instead of two round trips per record + The values are fetched in batches with a single HMGET call per batch, instead of several round trips per record as the default implementation does. The batches are bounded by the record sizes known from the metadata, so a store with large values does not load too many of them at once. """ diff --git a/src/crawlee/storage_clients/_utils.py b/src/crawlee/storage_clients/_utils.py index f02e285fd7..9fb65ba9ad 100644 --- a/src/crawlee/storage_clients/_utils.py +++ b/src/crawlee/storage_clients/_utils.py @@ -16,9 +16,9 @@ def batch_records_by_size( ) -> Iterator[list[KeyValueStoreRecordMetadata]]: """Group record metadata into batches bounded by record count and by total record size. - Used by storage clients that read the values of several records at once, so that a store with large values does - not load too many of them into memory in a single read. A record whose size alone exceeds `max_bytes` is still - yielded, but alone in its batch. A record with unknown size counts as empty. + Each batch is meant to be read in a single call, so the size bound keeps a store with large values from loading + too many of them into memory at once. A record whose size alone exceeds `max_bytes` is still yielded, but alone in + its batch. A record with unknown size counts as empty. """ batch: list[KeyValueStoreRecordMetadata] = [] batch_size = 0 diff --git a/src/crawlee/storages/_key_value_store.py b/src/crawlee/storages/_key_value_store.py index dcf6e7e73b..ae5be7e1b5 100644 --- a/src/crawlee/storages/_key_value_store.py +++ b/src/crawlee/storages/_key_value_store.py @@ -220,9 +220,9 @@ async def iterate_values( ) -> AsyncIterator[Any]: """Iterate over the values of the existing records in the KVS. - The records are fetched lazily as the iteration advances, so only a single value is held in memory at - a time. On remote backends this means one request per record on top of the paginated key listing, unless - the storage client provides a more efficient implementation. + The records are fetched lazily as the iteration advances, so only a bounded number of values is held in + memory at a time. By default this means one request per record on top of the paginated key listing, unless the + storage client reads the values in bounded batches instead. Args: exclusive_start_key: Key to start the iteration from. @@ -241,9 +241,9 @@ async def iterate_entries( ) -> AsyncIterator[tuple[str, Any]]: """Iterate over the existing records in the KVS as `(key, value)` pairs. - The records are fetched lazily as the iteration advances, so only a single value is held in memory at - a time. On remote backends this means one request per record on top of the paginated key listing, unless - the storage client provides a more efficient implementation. A record deleted while the iteration is in + The records are fetched lazily as the iteration advances, so only a bounded number of values is held in + memory at a time. By default this means one request per record on top of the paginated key listing, unless the + storage client reads the values in bounded batches instead. A record deleted while the iteration is in progress may or may not be yielded, depending on whether its value was already read. Args: From eb6a2f3a8bd5c9fa4c98f5e2f919bf3ce5d30f88 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 6 Oct 2026 10:37:00 +0200 Subject: [PATCH 15/22] test: cover SQL iterate_entries limit ending inside a later page --- .../storage_clients/_sql/test_sql_kvs_client.py | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py index 5aea03570a..4ff11de4a8 100644 --- a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py +++ b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py @@ -374,6 +374,20 @@ async def test_iterate_entries_with_exclusive_start_key_and_limit(kvs_client: Sq ] +async def test_iterate_entries_limit_spans_multiple_pages( + kvs_client: SqlKeyValueStoreClient, fetched_batches: list[list[str]] +) -> None: + """Test that `iterate_entries` stops at `limit` when it falls in the middle of a later metadata page.""" + for i in range(6): + await kvs_client.set_value(key=f'key{i}', value=f'value{i}') + + with patch.object(type(kvs_client), '_ITERATE_ENTRIES_BATCH_MAX_KEYS', 2): + records = [record async for record in kvs_client.iterate_entries(limit=3)] + + assert [(record.key, record.value) for record in records] == [(f'key{i}', f'value{i}') for i in range(3)] + assert fetched_batches == [['key0', 'key1'], ['key2']] + + async def test_iterate_entries_empty_store(kvs_client: SqlKeyValueStoreClient) -> None: records = [record async for record in kvs_client.iterate_entries()] From a967540faae8195e87372b0e6a52080d35b3bdf2 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 6 Oct 2026 10:37:01 +0200 Subject: [PATCH 16/22] test: rename batching test helpers and add missing test docstrings --- .../_redis/test_redis_kvs_client.py | 1 + .../_sql/test_sql_kvs_client.py | 1 + .../test_batch_records_by_size.py | 25 +++++++++++-------- 3 files changed, 17 insertions(+), 10 deletions(-) diff --git a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py index af2039633c..585bc9d353 100644 --- a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py +++ b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py @@ -358,6 +358,7 @@ async def test_iterate_entries_skips_value_deleted_after_listing( async def test_iterate_entries_empty_store(kvs_client: RedisKeyValueStoreClient) -> None: + """Test that `iterate_entries` on an empty store yields nothing.""" records = [record async for record in kvs_client.iterate_entries()] assert records == [] diff --git a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py index 4ff11de4a8..68b4174d94 100644 --- a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py +++ b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py @@ -389,6 +389,7 @@ async def test_iterate_entries_limit_spans_multiple_pages( async def test_iterate_entries_empty_store(kvs_client: SqlKeyValueStoreClient) -> None: + """Test that `iterate_entries` on an empty store yields nothing.""" records = [record async for record in kvs_client.iterate_entries()] assert records == [] diff --git a/tests/unit/storage_clients/test_batch_records_by_size.py b/tests/unit/storage_clients/test_batch_records_by_size.py index 2d546e8363..00f224ff13 100644 --- a/tests/unit/storage_clients/test_batch_records_by_size.py +++ b/tests/unit/storage_clients/test_batch_records_by_size.py @@ -4,45 +4,50 @@ from crawlee.storage_clients.models import KeyValueStoreRecordMetadata -def _record(key: str, size: int | None) -> KeyValueStoreRecordMetadata: +def make_record(key: str, size: int | None) -> KeyValueStoreRecordMetadata: return KeyValueStoreRecordMetadata(key=key, content_type='application/octet-stream', size=size) -def _keys(batches: list[list[KeyValueStoreRecordMetadata]]) -> list[list[str]]: +def batch_keys(batches: list[list[KeyValueStoreRecordMetadata]]) -> list[list[str]]: return [[record.key for record in batch] for batch in batches] def test_batches_are_bounded_by_record_count() -> None: - records = [_record(f'k{i}', 1) for i in range(5)] + """A new batch starts once the current one reaches `max_records`.""" + records = [make_record(f'k{i}', 1) for i in range(5)] batches = list(batch_records_by_size(records, max_records=2, max_bytes=1000)) - assert _keys(batches) == [['k0', 'k1'], ['k2', 'k3'], ['k4']] + assert batch_keys(batches) == [['k0', 'k1'], ['k2', 'k3'], ['k4']] def test_batches_are_bounded_by_total_size() -> None: - records = [_record('a', 4), _record('b', 4), _record('c', 4), _record('d', 4)] + """A new batch starts when the next record would push the total size over `max_bytes`.""" + records = [make_record('a', 4), make_record('b', 4), make_record('c', 4), make_record('d', 4)] batches = list(batch_records_by_size(records, max_records=100, max_bytes=10)) - assert _keys(batches) == [['a', 'b'], ['c', 'd']] + assert batch_keys(batches) == [['a', 'b'], ['c', 'd']] def test_oversized_record_is_alone_in_its_batch() -> None: - records = [_record('small1', 2), _record('large', 100), _record('small2', 2), _record('small3', 2)] + """A record larger than `max_bytes` is yielded alone in its batch.""" + records = [make_record('small1', 2), make_record('large', 100), make_record('small2', 2), make_record('small3', 2)] batches = list(batch_records_by_size(records, max_records=100, max_bytes=10)) - assert _keys(batches) == [['small1'], ['large'], ['small2', 'small3']] + assert batch_keys(batches) == [['small1'], ['large'], ['small2', 'small3']] def test_unknown_size_counts_as_empty() -> None: - records = [_record('a', None), _record('b', None), _record('c', 10)] + """A record with unknown size does not count toward `max_bytes`.""" + records = [make_record('a', None), make_record('b', None), make_record('c', 10)] batches = list(batch_records_by_size(records, max_records=100, max_bytes=10)) - assert _keys(batches) == [['a', 'b', 'c']] + assert batch_keys(batches) == [['a', 'b', 'c']] def test_no_records_yield_no_batches() -> None: + """An empty input yields no batches.""" assert list(batch_records_by_size([], max_records=100, max_bytes=10)) == [] From db0339a54854f22a23f313e46c424a23378774b7 Mon Sep 17 00:00:00 2001 From: Josef Prochazka Date: Tue, 6 Oct 2026 08:51:10 +0000 Subject: [PATCH 17/22] feat: make KeyValueStore async iteration yield keys like a dict `async for key in kvs` now yields the keys, matching how iterating over a `dict` behaves. Values and `(key, value)` pairs stay available through `iterate_values` and `iterate_entries`. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Wg6jZgyQp7XkvVvxuueJyV --- src/crawlee/storages/_key_value_store.py | 16 +++++++++------- tests/unit/storages/test_key_value_store.py | 6 +++--- 2 files changed, 12 insertions(+), 10 deletions(-) diff --git a/src/crawlee/storages/_key_value_store.py b/src/crawlee/storages/_key_value_store.py index ae5be7e1b5..4a7b310bd0 100644 --- a/src/crawlee/storages/_key_value_store.py +++ b/src/crawlee/storages/_key_value_store.py @@ -256,20 +256,22 @@ async def iterate_entries( async for record in self._client.iterate_entries(exclusive_start_key=exclusive_start_key, limit=limit): yield record.key, record.value - def __aiter__(self) -> AsyncIterator[tuple[str, Any]]: - """Iterate over all records in the KVS as `(key, value)` pairs. + async def __aiter__(self) -> AsyncIterator[str]: + """Iterate over all keys in the KVS. - Allows using the key-value store directly in an `async for` loop. It is equivalent to calling - `iterate_entries` with the default arguments. + Allows using the key-value store directly in an `async for` loop, which yields the keys like iterating + over a `dict` does. Use `iterate_keys` for the key metadata, `iterate_values` for the values, or + `iterate_entries` for `(key, value)` pairs. ### Usage ```python - async for key, value in kvs: - print(key, value) + async for key in kvs: + print(key) ``` """ - return self.iterate_entries() + async for metadata in self.iterate_keys(): + yield metadata.key async def list_keys( self, diff --git a/tests/unit/storages/test_key_value_store.py b/tests/unit/storages/test_key_value_store.py index 3a8aa0d9c0..41ef540033 100644 --- a/tests/unit/storages/test_key_value_store.py +++ b/tests/unit/storages/test_key_value_store.py @@ -359,13 +359,13 @@ async def get_value(self, *, key: str) -> KeyValueStoreRecord | None: async def test_async_iteration(kvs: KeyValueStore) -> None: - """Test that the key-value store can be used directly in an `async for` loop, yielding (key, value) pairs.""" + """Test that the key-value store can be used directly in an `async for` loop, yielding keys like a dict.""" await kvs.set_value('key1', 'value1') await kvs.set_value('key2', 'value2') - collected_entries = {key: value async for key, value in kvs} + collected_keys = [key async for key in kvs] - assert collected_entries == {'key1': 'value1', 'key2': 'value2'} + assert sorted(collected_keys) == ['key1', 'key2'] async def test_drop( From 0b9af818e00f396c592e013a46978cb955983e3c Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 6 Oct 2026 11:48:44 +0200 Subject: [PATCH 18/22] refactor: return the key iterator from KeyValueStore.__aiter__ like Dataset does --- src/crawlee/storages/_key_value_store.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/src/crawlee/storages/_key_value_store.py b/src/crawlee/storages/_key_value_store.py index 4a7b310bd0..b6b59d76d5 100644 --- a/src/crawlee/storages/_key_value_store.py +++ b/src/crawlee/storages/_key_value_store.py @@ -256,7 +256,7 @@ async def iterate_entries( async for record in self._client.iterate_entries(exclusive_start_key=exclusive_start_key, limit=limit): yield record.key, record.value - async def __aiter__(self) -> AsyncIterator[str]: + def __aiter__(self) -> AsyncIterator[str]: """Iterate over all keys in the KVS. Allows using the key-value store directly in an `async for` loop, which yields the keys like iterating @@ -270,8 +270,7 @@ async def __aiter__(self) -> AsyncIterator[str]: print(key) ``` """ - async for metadata in self.iterate_keys(): - yield metadata.key + return (metadata.key async for metadata in self.iterate_keys()) async def list_keys( self, From 4905199ccd2dba097cdd61c84406126528129ca1 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 6 Oct 2026 12:06:26 +0200 Subject: [PATCH 19/22] docs: fix RedisStorageClient and KVS iteration docstrings --- src/crawlee/storage_clients/_redis/_storage_client.py | 8 ++++---- src/crawlee/storages/_key_value_store.py | 8 ++++---- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/src/crawlee/storage_clients/_redis/_storage_client.py b/src/crawlee/storage_clients/_redis/_storage_client.py index 3450df1390..bc8c8759e0 100644 --- a/src/crawlee/storage_clients/_redis/_storage_client.py +++ b/src/crawlee/storage_clients/_redis/_storage_client.py @@ -23,10 +23,10 @@ class RedisStorageClient(StorageClient): to a Redis database v8.0+. Each storage type uses Redis-specific data structures and key patterns for efficient storage and retrieval. - The client accepts either a Redis connection string or a pre-configured Redis client instance. Key-value stores - hold binary values, so their client needs a Redis client that returns raw bytes, which is the default; a client - created with `decode_responses=True` works for datasets and request queues only. - Exactly one of these parameters must be provided during initialization. + The client accepts either a Redis connection string or a pre-configured Redis client instance. Exactly one of these + parameters must be provided during initialization. Key-value stores hold binary values, so their client needs + a Redis client that returns raw bytes, which is the default; a client created with `decode_responses=True` works + for datasets and request queues only. Storage types use the following Redis data structures: - **Datasets**: Redis JSON arrays for item storage with metadata in JSON objects diff --git a/src/crawlee/storages/_key_value_store.py b/src/crawlee/storages/_key_value_store.py index b6b59d76d5..1eb82087e1 100644 --- a/src/crawlee/storages/_key_value_store.py +++ b/src/crawlee/storages/_key_value_store.py @@ -221,8 +221,8 @@ async def iterate_values( """Iterate over the values of the existing records in the KVS. The records are fetched lazily as the iteration advances, so only a bounded number of values is held in - memory at a time. By default this means one request per record on top of the paginated key listing, unless the - storage client reads the values in bounded batches instead. + memory at a time. By default this means one request per record on top of the key listing, unless the storage + client reads the values in bounded batches instead. Args: exclusive_start_key: Key to start the iteration from. @@ -242,8 +242,8 @@ async def iterate_entries( """Iterate over the existing records in the KVS as `(key, value)` pairs. The records are fetched lazily as the iteration advances, so only a bounded number of values is held in - memory at a time. By default this means one request per record on top of the paginated key listing, unless the - storage client reads the values in bounded batches instead. A record deleted while the iteration is in + memory at a time. By default this means one request per record on top of the key listing, unless the storage + client reads the values in bounded batches instead. A record deleted while the iteration is in progress may or may not be yielded, depending on whether its value was already read. Args: From 1c66c3b9b29babdfd3de5bfa447f98f992fecf2f Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 6 Oct 2026 12:06:27 +0200 Subject: [PATCH 20/22] test: record Redis and SQL batch calls with wrapping mocks --- .../_redis/test_redis_kvs_client.py | 31 ++++++++----------- .../_sql/test_sql_kvs_client.py | 29 ++++++++--------- 2 files changed, 26 insertions(+), 34 deletions(-) diff --git a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py index 585bc9d353..69dbcc5d6d 100644 --- a/tests/unit/storage_clients/_redis/test_redis_kvs_client.py +++ b/tests/unit/storage_clients/_redis/test_redis_kvs_client.py @@ -2,7 +2,7 @@ import asyncio import json -from typing import TYPE_CHECKING, Any +from typing import TYPE_CHECKING from unittest.mock import AsyncMock, MagicMock, patch import pytest @@ -259,22 +259,17 @@ async def test_set_value_does_not_retry_on_unexpected_exception(kvs_client: Redi @pytest.fixture -def hmget_calls(kvs_client: RedisKeyValueStoreClient) -> Iterator[list[list[str]]]: - """Record the keys of every Redis `hmget` call made through the client, while still performing the call.""" - calls: list[list[str]] = [] - original_hmget = kvs_client.redis.hmget +def hmget(kvs_client: RedisKeyValueStoreClient) -> Iterator[MagicMock]: + """Wrap the Redis `hmget` of the client in a mock that records the calls while still performing them.""" + with patch.object(kvs_client.redis, 'hmget', wraps=kvs_client.redis.hmget) as mock: + yield mock - def recording_hmget(name: str, keys: list[str], *args: str) -> Any: - calls.append(list(keys)) - return original_hmget(name, keys, *args) - with patch.object(kvs_client.redis, 'hmget', side_effect=recording_hmget): - yield calls +def hmget_keys(hmget: MagicMock) -> list[list[str]]: + return [list(call.args[1]) for call in hmget.call_args_list] -async def test_iterate_entries_reads_values_in_batches( - kvs_client: RedisKeyValueStoreClient, hmget_calls: list[list[str]] -) -> None: +async def test_iterate_entries_reads_values_in_batches(kvs_client: RedisKeyValueStoreClient, hmget: MagicMock) -> None: """Test that `iterate_entries` fetches values with batched HMGET calls instead of `get_value` per key.""" await kvs_client.set_value(key='a-json', value={'nested': [1, 2]}) await kvs_client.set_value(key='b-text', value='plain text') @@ -291,11 +286,11 @@ async def test_iterate_entries_reads_values_in_batches( assert records[2].content_type == 'application/octet-stream' # All values fit in a single batch. - assert hmget_calls == [['a-json', 'b-text', 'c-bytes', 'd-none']] + assert hmget_keys(hmget) == [['a-json', 'b-text', 'c-bytes', 'd-none']] async def test_iterate_entries_batches_are_bounded_by_key_count( - kvs_client: RedisKeyValueStoreClient, hmget_calls: list[list[str]] + kvs_client: RedisKeyValueStoreClient, hmget: MagicMock ) -> None: """Test that `iterate_entries` splits the HMGET calls when a batch reaches the maximum number of keys.""" for i in range(5): @@ -305,11 +300,11 @@ async def test_iterate_entries_batches_are_bounded_by_key_count( records = [record async for record in kvs_client.iterate_entries()] assert [(record.key, record.value) for record in records] == [(f'key{i}', f'value{i}') for i in range(5)] - assert hmget_calls == [['key0', 'key1'], ['key2', 'key3'], ['key4']] + assert hmget_keys(hmget) == [['key0', 'key1'], ['key2', 'key3'], ['key4']] async def test_iterate_entries_batches_are_bounded_by_size( - kvs_client: RedisKeyValueStoreClient, hmget_calls: list[list[str]] + kvs_client: RedisKeyValueStoreClient, hmget: MagicMock ) -> None: """Test that `iterate_entries` splits the HMGET calls by the record sizes known from the metadata. @@ -324,7 +319,7 @@ async def test_iterate_entries_batches_are_bounded_by_size( records = [record async for record in kvs_client.iterate_entries()] assert [record.key for record in records] == ['large', 'small1', 'small2', 'small3'] - assert hmget_calls == [['large'], ['small1', 'small2', 'small3']] + assert hmget_keys(hmget) == [['large'], ['small1', 'small2', 'small3']] async def test_iterate_entries_with_exclusive_start_key_and_limit(kvs_client: RedisKeyValueStoreClient) -> None: diff --git a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py index 68b4174d94..5c2c7a3300 100644 --- a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py +++ b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py @@ -13,7 +13,7 @@ from crawlee.configuration import Configuration from crawlee.storage_clients import SqlStorageClient from crawlee.storage_clients._sql._db_models import KeyValueStoreMetadataDb, KeyValueStoreRecordDb -from crawlee.storage_clients.models import KeyValueStoreMetadata, KeyValueStoreRecord, KeyValueStoreRecordMetadata +from crawlee.storage_clients.models import KeyValueStoreMetadata if TYPE_CHECKING: from collections.abc import AsyncGenerator, Iterator @@ -328,17 +328,14 @@ async def test_set_value_does_not_retry_on_unexpected_exception(kvs_client: SqlK @pytest.fixture -def fetched_batches(kvs_client: SqlKeyValueStoreClient) -> Iterator[list[list[str]]]: - """Record the keys of every value batch `iterate_entries` fetches, while still performing the fetch.""" - calls: list[list[str]] = [] - original_fetch = kvs_client._fetch_records +def fetch_records(kvs_client: SqlKeyValueStoreClient) -> Iterator[AsyncMock]: + """Wrap `_fetch_records` of the client in a mock that records the value batches while still fetching them.""" + with patch.object(kvs_client, '_fetch_records', wraps=kvs_client._fetch_records) as mock: + yield mock - async def recording_fetch(batch: list[KeyValueStoreRecordMetadata]) -> list[KeyValueStoreRecord]: - calls.append([item.key for item in batch]) - return await original_fetch(batch) - with patch.object(kvs_client, '_fetch_records', side_effect=recording_fetch): - yield calls +def fetched_keys(fetch_records: AsyncMock) -> list[list[str]]: + return [[item.key for item in call.args[0]] for call in fetch_records.await_args_list] async def test_iterate_entries_reads_values_in_batches(kvs_client: SqlKeyValueStoreClient) -> None: @@ -375,7 +372,7 @@ async def test_iterate_entries_with_exclusive_start_key_and_limit(kvs_client: Sq async def test_iterate_entries_limit_spans_multiple_pages( - kvs_client: SqlKeyValueStoreClient, fetched_batches: list[list[str]] + kvs_client: SqlKeyValueStoreClient, fetch_records: AsyncMock ) -> None: """Test that `iterate_entries` stops at `limit` when it falls in the middle of a later metadata page.""" for i in range(6): @@ -385,7 +382,7 @@ async def test_iterate_entries_limit_spans_multiple_pages( records = [record async for record in kvs_client.iterate_entries(limit=3)] assert [(record.key, record.value) for record in records] == [(f'key{i}', f'value{i}') for i in range(3)] - assert fetched_batches == [['key0', 'key1'], ['key2']] + assert fetched_keys(fetch_records) == [['key0', 'key1'], ['key2']] async def test_iterate_entries_empty_store(kvs_client: SqlKeyValueStoreClient) -> None: @@ -396,7 +393,7 @@ async def test_iterate_entries_empty_store(kvs_client: SqlKeyValueStoreClient) - async def test_iterate_entries_batches_are_bounded_by_key_count( - kvs_client: SqlKeyValueStoreClient, fetched_batches: list[list[str]] + kvs_client: SqlKeyValueStoreClient, fetch_records: AsyncMock ) -> None: """Test that `iterate_entries` splits the value queries when a batch reaches the maximum number of keys.""" for i in range(5): @@ -406,11 +403,11 @@ async def test_iterate_entries_batches_are_bounded_by_key_count( records = [record async for record in kvs_client.iterate_entries()] assert [(record.key, record.value) for record in records] == [(f'key{i}', f'value{i}') for i in range(5)] - assert fetched_batches == [['key0', 'key1'], ['key2', 'key3'], ['key4']] + assert fetched_keys(fetch_records) == [['key0', 'key1'], ['key2', 'key3'], ['key4']] async def test_iterate_entries_batches_are_bounded_by_size( - kvs_client: SqlKeyValueStoreClient, fetched_batches: list[list[str]] + kvs_client: SqlKeyValueStoreClient, fetch_records: AsyncMock ) -> None: """Test that `iterate_entries` splits the value queries by the record sizes. @@ -425,4 +422,4 @@ async def test_iterate_entries_batches_are_bounded_by_size( records = [record async for record in kvs_client.iterate_entries()] assert [record.key for record in records] == ['large', 'small1', 'small2', 'small3'] - assert fetched_batches == [['large'], ['small1', 'small2', 'small3']] + assert fetched_keys(fetch_records) == [['large'], ['small1', 'small2', 'small3']] From 2c3dd78e20a177e3a25f0188d564991d76f61a6f Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 6 Oct 2026 12:06:28 +0200 Subject: [PATCH 21/22] test: cover records deleted mid-iteration in memory and SQL KVS clients --- .../_memory/test_memory_kvs_client.py | 14 ++++++++++++++ .../storage_clients/_sql/test_sql_kvs_client.py | 12 ++++++++++++ 2 files changed, 26 insertions(+) diff --git a/tests/unit/storage_clients/_memory/test_memory_kvs_client.py b/tests/unit/storage_clients/_memory/test_memory_kvs_client.py index 4dfc44085e..4c7fce8058 100644 --- a/tests/unit/storage_clients/_memory/test_memory_kvs_client.py +++ b/tests/unit/storage_clients/_memory/test_memory_kvs_client.py @@ -78,3 +78,17 @@ async def test_memory_metadata_updates(kvs_client: MemoryKeyValueStoreClient) -> assert metadata.created_at == initial_created assert metadata.modified_at > initial_modified assert metadata.accessed_at > accessed_after_read + + +async def test_iterate_keys_skips_record_deleted_during_iteration(kvs_client: MemoryKeyValueStoreClient) -> None: + """Test that a record deleted while `iterate_keys` is suspended is skipped.""" + await kvs_client.set_value(key='key1', value='a') + await kvs_client.set_value(key='key2', value='b') + + collected_keys = [] + async for metadata in kvs_client.iterate_keys(): + if metadata.key == 'key1': + await kvs_client.delete_value(key='key2') + collected_keys.append(metadata.key) + + assert collected_keys == ['key1'] diff --git a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py index 5c2c7a3300..cfc2722db8 100644 --- a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py +++ b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py @@ -385,6 +385,18 @@ async def test_iterate_entries_limit_spans_multiple_pages( assert fetched_keys(fetch_records) == [['key0', 'key1'], ['key2']] +async def test_iterate_entries_skips_record_deleted_after_listing(kvs_client: SqlKeyValueStoreClient) -> None: + """Test that a record deleted between the metadata listing and its value fetch is skipped.""" + await kvs_client.set_value(key='kept', value='a') + await kvs_client.set_value(key='removed', value='b') + batch = [item async for item in kvs_client.iterate_keys()] + await kvs_client.delete_value(key='removed') + + records = await kvs_client._fetch_records(batch) + + assert [(record.key, record.value) for record in records] == [('kept', 'a')] + + async def test_iterate_entries_empty_store(kvs_client: SqlKeyValueStoreClient) -> None: """Test that `iterate_entries` on an empty store yields nothing.""" records = [record async for record in kvs_client.iterate_entries()] From b4c21ae647e66da6e260dde9215c1fac875dcccc Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Thu, 8 Oct 2026 09:17:42 +0200 Subject: [PATCH 22/22] refactor: tidy iterate_entries code, docstrings and test names --- .../storage_clients/_base/_key_value_store_client.py | 5 ++--- src/crawlee/storage_clients/_sql/_key_value_store_client.py | 2 +- tests/unit/storage_clients/_sql/test_sql_kvs_client.py | 2 +- tests/unit/storages/test_key_value_store.py | 6 +----- 4 files changed, 5 insertions(+), 10 deletions(-) diff --git a/src/crawlee/storage_clients/_base/_key_value_store_client.py b/src/crawlee/storage_clients/_base/_key_value_store_client.py index 80e73f7af5..53d49059eb 100644 --- a/src/crawlee/storage_clients/_base/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_base/_key_value_store_client.py @@ -97,9 +97,8 @@ async def iterate_entries( """ async for metadata in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit): record = await self.get_value(key=metadata.key) - if record is None: - continue - yield record + if record is not None: + yield record @abstractmethod async def get_public_url(self, *, key: str) -> str: diff --git a/src/crawlee/storage_clients/_sql/_key_value_store_client.py b/src/crawlee/storage_clients/_sql/_key_value_store_client.py index d2d97d1dc6..368379439a 100644 --- a/src/crawlee/storage_clients/_sql/_key_value_store_client.py +++ b/src/crawlee/storage_clients/_sql/_key_value_store_client.py @@ -61,7 +61,7 @@ class SqlKeyValueStoreClient(KeyValueStoreClient, SqlClientMixin): """Default dataset name used when no name is provided.""" _ITERATE_ENTRIES_BATCH_MAX_KEYS = 100 - """Maximum number of records read with a single query in `iterate_entries`.""" + """Maximum number of records listed or read with a single query in `iterate_entries`.""" _ITERATE_ENTRIES_BATCH_MAX_BYTES = 8 * 1024 * 1024 """Maximum total size of the records read with a single query in `iterate_entries`. diff --git a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py index cfc2722db8..1b31165964 100644 --- a/tests/unit/storage_clients/_sql/test_sql_kvs_client.py +++ b/tests/unit/storage_clients/_sql/test_sql_kvs_client.py @@ -385,7 +385,7 @@ async def test_iterate_entries_limit_spans_multiple_pages( assert fetched_keys(fetch_records) == [['key0', 'key1'], ['key2']] -async def test_iterate_entries_skips_record_deleted_after_listing(kvs_client: SqlKeyValueStoreClient) -> None: +async def test_fetch_records_skips_record_deleted_after_listing(kvs_client: SqlKeyValueStoreClient) -> None: """Test that a record deleted between the metadata listing and its value fetch is skipped.""" await kvs_client.set_value(key='kept', value='a') await kvs_client.set_value(key='removed', value='b') diff --git a/tests/unit/storages/test_key_value_store.py b/tests/unit/storages/test_key_value_store.py index 41ef540033..1faadce748 100644 --- a/tests/unit/storages/test_key_value_store.py +++ b/tests/unit/storages/test_key_value_store.py @@ -318,11 +318,7 @@ async def test_iterate_entries_empty_kvs(kvs: KeyValueStore) -> None: async def test_iterate_entries_uses_storage_client_implementation() -> None: - """Test that `iterate_entries` and `iterate_values` go through the storage client's `iterate_entries`. - - Storage clients can override the default key-by-key implementation with a more efficient one, so the frontend - must delegate to the client instead of combining `iterate_keys` and `get_value` itself. - """ + """Test that `iterate_entries` and `iterate_values` delegate to the storage client's `iterate_entries`.""" class OptimizedKeyValueStoreClient(MemoryKeyValueStoreClient): async def iterate_entries(