Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
193cc2e
feat: add direct async iteration on Dataset and KeyValueStore
Pijukatel Oct 1, 2026
2deace2
feat: make KeyValueStore record iteration overridable by storage clients
Pijukatel Oct 1, 2026
3c03eae
perf: read key-value store entries in bulk in the SQL and Redis clients
Pijukatel Oct 1, 2026
125bba2
test: drop the test for records deleted during key-value store iteration
Pijukatel Oct 2, 2026
668c170
refactor: decode None values in one place in the Redis key-value stor…
Pijukatel Oct 2, 2026
0e255a5
docs: describe the bulk iterate_entries overrides in docstrings
Pijukatel Oct 2, 2026
6c0cc3f
refactor: reject a decoding Redis client instead of suppressing the t…
Pijukatel Oct 2, 2026
95a5154
refactor: narrow Redis replies without a cast
Pijukatel Oct 2, 2026
a07f81a
docs: scope the raw-bytes requirement to the Redis key-value store cl…
Pijukatel Oct 2, 2026
1f32875
fix: bound memory of the SQL key-value store iteration with keyset pages
Pijukatel Oct 2, 2026
f809fe7
fix: retry the batched Redis value read and test the batching helper
Pijukatel Oct 2, 2026
c0c334e
test: give the batching helper test module a unique basename
Pijukatel Oct 2, 2026
f7db136
fix: skip Redis records deleted mid-iteration without a missing-value…
vdusek Oct 6, 2026
5918899
docs: correct memory and round-trip claims in KVS iteration docstrings
vdusek Oct 6, 2026
eb6a2f3
test: cover SQL iterate_entries limit ending inside a later page
vdusek Oct 6, 2026
a967540
test: rename batching test helpers and add missing test docstrings
vdusek Oct 6, 2026
db0339a
feat: make KeyValueStore async iteration yield keys like a dict
Pijukatel Oct 6, 2026
0b9af81
refactor: return the key iterator from KeyValueStore.__aiter__ like D…
vdusek Oct 6, 2026
4905199
docs: fix RedisStorageClient and KVS iteration docstrings
vdusek Oct 6, 2026
1c66c3b
test: record Redis and SQL batch calls with wrapping mocks
vdusek Oct 6, 2026
2c3dd78
test: cover records deleted mid-iteration in memory and SQL KVS clients
vdusek Oct 6, 2026
b4c21ae
refactor: tidy iterate_entries code, docstrings and test names
vdusek Oct 8, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions src/crawlee/storage_clients/_base/_key_value_store_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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. 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)
if record is not None:
yield record

@abstractmethod
async def get_public_url(self, *, key: str) -> str:
"""Get the public URL for the given key.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
75 changes: 71 additions & 4 deletions src/crawlee/storage_clients/_redis/_key_value_store_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,11 @@
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
from ._utils import await_redis_response
from ._utils import await_redis_response, expect_bytes

if TYPE_CHECKING:
from collections.abc import AsyncIterator
Expand All @@ -40,6 +41,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'
Expand All @@ -51,6 +55,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.

Expand Down Expand Up @@ -178,15 +191,30 @@ 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)

@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

# 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):
Expand Down Expand Up @@ -252,6 +280,45 @@ 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]:
"""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 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.
"""
metadata_items = [
item async for item in self.iterate_keys(exclusive_start_key=exclusive_start_key, limit=limit)
]
Comment thread
Pijukatel marked this conversation as resolved.

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,
):
for record in await self._fetch_records(batch):
yield record

@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))]

# 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]

@override
async def get_public_url(self, *, key: str) -> str:
raise NotImplementedError('Public URLs are not supported for memory key-value stores.')
Expand Down
6 changes: 4 additions & 2 deletions src/crawlee/storage_clients/_redis/_storage_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +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.
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
Expand Down
18 changes: 18 additions & 0 deletions src/crawlee/storage_clients/_redis/_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,24 @@ async def await_redis_response(response: Awaitable[T] | T) -> T:
return response


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 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 key-value store client '
'requires a Redis client created without `decode_responses=True`.'
)
return 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
Expand Down
129 changes: 118 additions & 11 deletions src/crawlee/storage_clients/_sql/_key_value_store_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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 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`.

A single record larger than this is still read, but alone in its batch.
"""

_METADATA_TABLE = KeyValueStoreMetadataDb
"""SQLAlchemy model for key-value store metadata."""

Expand Down Expand Up @@ -202,21 +212,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:
Expand All @@ -226,12 +248,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
Expand Down Expand Up @@ -282,6 +299,96 @@ 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]:
"""Iterate over all the existing records in the key-value store, including their values.

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,
self._ITEM_TABLE.content_type,
self._ITEM_TABLE.size,
self._ITEM_TABLE.value,
)
.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)
)

async with self.get_session(with_simple_commit=True) as session:
result = await session.execute(stmt)
rows = result.all()

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
async def record_exists(self, *, key: str) -> bool:
Expand Down
Loading
Loading