Repository navigation
Expand file tree
/
Copy path_storage_client.py
More file actions
149 lines (123 loc) · 5.39 KB
/
Copy path_storage_client.py
File metadata and controls
149 lines (123 loc) · 5.39 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
from __future__ import annotations
import warnings
from typing import Literal
from redis.asyncio import Redis
from typing_extensions import override
from crawlee._utils.docs import docs_group
from crawlee.configuration import Configuration
from crawlee.storage_clients._base import StorageClient
from ._dataset_client import RedisDatasetClient
from ._key_value_store_client import RedisKeyValueStoreClient
from ._request_queue_client import RedisRequestQueueClient
@docs_group('Storage clients')
class RedisStorageClient(StorageClient):
"""Redis implementation of the storage client.
This storage client provides access to datasets, key-value stores, and request queues that persist data
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.
Storage types use the following Redis data structures:
- **Datasets**: Redis JSON arrays for item storage with metadata in JSON objects
- **Key-value stores**: Redis hashes for key-value pairs with separate metadata storage
- **Request queues**: Redis lists for FIFO queuing, hashes for request data and in-progress tracking,
and Bloom filters for request deduplication
Warning:
This is an experimental feature. The behavior and interface may change in future versions.
"""
def __init__(
self,
*,
connection_string: str | None = None,
redis: Redis | None = None,
queue_dedup_strategy: Literal['default', 'bloom'] = 'default',
queue_bloom_error_rate: float = 1e-7,
) -> None:
"""Initialize the Redis storage client.
Args:
connection_string: Redis connection string (e.g., "redis://localhost:6379").
Supports standard Redis URL format with optional database selection.
redis: Pre-configured Redis client instance.
queue_dedup_strategy: Strategy for request queue deduplication. Options are:
- 'default': Uses Redis sets for exact deduplication.
- 'bloom': Uses Redis Bloom filters for probabilistic deduplication with lower memory usage. When using
this approach, approximately 1 in 1e-7 requests will be falsely considered duplicate.
queue_bloom_error_rate: Desired false positive rate for Bloom filter deduplication. Only relevant if
`queue_dedup_strategy` is set to 'bloom'.
"""
if redis is None and connection_string is None:
raise ValueError('Either redis or connection_string must be provided.')
if redis is not None and connection_string is not None:
raise ValueError('Either redis or connection_string must be provided, not both.')
if isinstance(redis, Redis) and connection_string is None:
self._redis = redis
if isinstance(connection_string, str) and redis is None:
self._redis = Redis.from_url(connection_string)
self._redis: Redis # to help type checker
self._queue_dedup_strategy = queue_dedup_strategy
self._queue_bloom_error_rate = queue_bloom_error_rate
# Call the notification only once
warnings.warn(
(
'RedisStorageClient is experimental and its API, behavior, and key structure may change in future '
'releases.'
),
category=UserWarning,
stacklevel=2,
)
@override
async def create_dataset_client(
self,
*,
id: str | None = None,
name: str | None = None,
alias: str | None = None,
configuration: Configuration | None = None,
) -> RedisDatasetClient:
configuration = configuration or Configuration.get_global_configuration()
client = await RedisDatasetClient.open(
id=id,
name=name,
alias=alias,
redis=self._redis,
)
await self._purge_if_needed(client, configuration)
return client
@override
async def create_kvs_client(
self,
*,
id: str | None = None,
name: str | None = None,
alias: str | None = None,
configuration: Configuration | None = None,
) -> RedisKeyValueStoreClient:
configuration = configuration or Configuration.get_global_configuration()
client = await RedisKeyValueStoreClient.open(
id=id,
name=name,
alias=alias,
redis=self._redis,
)
await self._purge_if_needed(client, configuration)
return client
@override
async def create_rq_client(
self,
*,
id: str | None = None,
name: str | None = None,
alias: str | None = None,
configuration: Configuration | None = None,
) -> RedisRequestQueueClient:
configuration = configuration or Configuration.get_global_configuration()
client = await RedisRequestQueueClient.open(
id=id,
name=name,
alias=alias,
redis=self._redis,
dedup_strategy=self._queue_dedup_strategy,
bloom_error_rate=self._queue_bloom_error_rate,
)
await self._purge_if_needed(client, configuration)
return client