变更
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
from .connection import get_redis, get_redis_from_settings # NOQA
|
||||
|
||||
__author__ = "R Max Espinoza"
|
||||
__email__ = "hey at rmax.dev"
|
||||
__version__ = "0.9.1"
|
||||
@@ -0,0 +1,97 @@
|
||||
from scrapy.utils.misc import load_object
|
||||
|
||||
from . import defaults
|
||||
|
||||
# Shortcut maps 'setting name' -> 'parmater name'.
|
||||
SETTINGS_PARAMS_MAP = {
|
||||
"REDIS_URL": "url",
|
||||
"REDIS_HOST": "host",
|
||||
"REDIS_PORT": "port",
|
||||
"REDIS_DB": "db",
|
||||
"REDIS_ENCODING": "encoding",
|
||||
}
|
||||
|
||||
SETTINGS_PARAMS_MAP["REDIS_DECODE_RESPONSES"] = "decode_responses"
|
||||
|
||||
|
||||
def get_redis_from_settings(settings):
|
||||
"""Returns a redis client instance from given Scrapy settings object.
|
||||
|
||||
This function uses ``get_client`` to instantiate the client and uses
|
||||
``defaults.REDIS_PARAMS`` global as defaults values for the parameters. You
|
||||
can override them using the ``REDIS_PARAMS`` setting.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
settings : Settings
|
||||
A scrapy settings object. See the supported settings below.
|
||||
|
||||
Returns
|
||||
-------
|
||||
server
|
||||
Redis client instance.
|
||||
|
||||
Other Parameters
|
||||
----------------
|
||||
REDIS_URL : str, optional
|
||||
Server connection URL.
|
||||
REDIS_HOST : str, optional
|
||||
Server host.
|
||||
REDIS_PORT : str, optional
|
||||
Server port.
|
||||
REDIS_DB : int, optional
|
||||
Server database
|
||||
REDIS_ENCODING : str, optional
|
||||
Data encoding.
|
||||
REDIS_PARAMS : dict, optional
|
||||
Additional client parameters.
|
||||
|
||||
Python 3 Only
|
||||
----------------
|
||||
REDIS_DECODE_RESPONSES : bool, optional
|
||||
Sets the `decode_responses` kwarg in Redis cls ctor
|
||||
|
||||
"""
|
||||
params = defaults.REDIS_PARAMS.copy()
|
||||
params.update(settings.getdict("REDIS_PARAMS"))
|
||||
# XXX: Deprecate REDIS_* settings.
|
||||
for source, dest in SETTINGS_PARAMS_MAP.items():
|
||||
val = settings.get(source)
|
||||
if val:
|
||||
params[dest] = val
|
||||
|
||||
# Allow ``redis_cls`` to be a path to a class.
|
||||
if isinstance(params.get("redis_cls"), str):
|
||||
params["redis_cls"] = load_object(params["redis_cls"])
|
||||
|
||||
return get_redis(**params)
|
||||
|
||||
|
||||
# Backwards compatible alias.
|
||||
from_settings = get_redis_from_settings
|
||||
|
||||
|
||||
def get_redis(**kwargs):
|
||||
"""Returns a redis client instance.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
redis_cls : class, optional
|
||||
Defaults to ``redis.StrictRedis``.
|
||||
url : str, optional
|
||||
If given, ``redis_cls.from_url`` is used to instantiate the class.
|
||||
**kwargs
|
||||
Extra parameters to be passed to the ``redis_cls`` class.
|
||||
|
||||
Returns
|
||||
-------
|
||||
server
|
||||
Redis client instance.
|
||||
|
||||
"""
|
||||
redis_cls = kwargs.pop("redis_cls", defaults.REDIS_CLS)
|
||||
url = kwargs.pop("url", None)
|
||||
if url:
|
||||
return redis_cls.from_url(url, **kwargs)
|
||||
else:
|
||||
return redis_cls(**kwargs)
|
||||
@@ -0,0 +1,29 @@
|
||||
import redis
|
||||
|
||||
# For standalone use.
|
||||
DUPEFILTER_KEY = "dupefilter:%(timestamp)s"
|
||||
|
||||
PIPELINE_KEY = "%(spider)s:items"
|
||||
|
||||
STATS_KEY = "%(spider)s:stats"
|
||||
|
||||
REDIS_CLS = redis.StrictRedis
|
||||
REDIS_ENCODING = "utf-8"
|
||||
# Sane connection defaults.
|
||||
REDIS_PARAMS = {
|
||||
"socket_timeout": 30,
|
||||
"socket_connect_timeout": 30,
|
||||
"retry_on_timeout": True,
|
||||
"encoding": REDIS_ENCODING,
|
||||
}
|
||||
REDIS_CONCURRENT_REQUESTS = 16
|
||||
|
||||
SCHEDULER_QUEUE_KEY = "%(spider)s:requests"
|
||||
SCHEDULER_QUEUE_CLASS = "scrapy_redis.queue.PriorityQueue"
|
||||
SCHEDULER_DUPEFILTER_KEY = "%(spider)s:dupefilter"
|
||||
SCHEDULER_DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
|
||||
SCHEDULER_PERSIST = False
|
||||
START_URLS_KEY = "%(name)s:start_urls"
|
||||
START_URLS_AS_SET = False
|
||||
START_URLS_AS_ZSET = False
|
||||
MAX_IDLE_TIME = 0
|
||||
@@ -0,0 +1,169 @@
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import time
|
||||
|
||||
from scrapy.dupefilters import BaseDupeFilter
|
||||
from scrapy.utils.python import to_unicode
|
||||
from w3lib.url import canonicalize_url
|
||||
|
||||
from . import defaults
|
||||
from .connection import get_redis_from_settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# TODO: Rename class to RedisDupeFilter.
|
||||
class RFPDupeFilter(BaseDupeFilter):
|
||||
"""Redis-based request duplicates filter.
|
||||
|
||||
This class can also be used with default Scrapy's scheduler.
|
||||
|
||||
"""
|
||||
|
||||
logger = logger
|
||||
|
||||
def __init__(self, server, key, debug=False):
|
||||
"""Initialize the duplicates filter.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
server : redis.StrictRedis
|
||||
The redis server instance.
|
||||
key : str
|
||||
Redis key Where to store fingerprints.
|
||||
debug : bool, optional
|
||||
Whether to log filtered requests.
|
||||
|
||||
"""
|
||||
self.server = server
|
||||
self.key = key
|
||||
self.debug = debug
|
||||
self.logdupes = True
|
||||
|
||||
@classmethod
|
||||
def from_settings(cls, settings):
|
||||
"""Returns an instance from given settings.
|
||||
|
||||
This uses by default the key ``dupefilter:<timestamp>``. When using the
|
||||
``scrapy_redis.scheduler.Scheduler`` class, this method is not used as
|
||||
it needs to pass the spider name in the key.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
settings : scrapy.settings.Settings
|
||||
|
||||
Returns
|
||||
-------
|
||||
RFPDupeFilter
|
||||
A RFPDupeFilter instance.
|
||||
|
||||
|
||||
"""
|
||||
server = get_redis_from_settings(settings)
|
||||
# XXX: This creates one-time key. needed to support to use this
|
||||
# class as standalone dupefilter with scrapy's default scheduler
|
||||
# if scrapy passes spider on open() method this wouldn't be needed
|
||||
# TODO: Use SCRAPY_JOB env as default and fallback to timestamp.
|
||||
key = defaults.DUPEFILTER_KEY % {"timestamp": int(time.time())}
|
||||
debug = settings.getbool("DUPEFILTER_DEBUG")
|
||||
return cls(server, key=key, debug=debug)
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler):
|
||||
"""Returns instance from crawler.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
crawler : scrapy.crawler.Crawler
|
||||
|
||||
Returns
|
||||
-------
|
||||
RFPDupeFilter
|
||||
Instance of RFPDupeFilter.
|
||||
|
||||
"""
|
||||
return cls.from_settings(crawler.settings)
|
||||
|
||||
def request_seen(self, request):
|
||||
"""Returns True if request was already seen.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
request : scrapy.http.Request
|
||||
|
||||
Returns
|
||||
-------
|
||||
bool
|
||||
|
||||
"""
|
||||
fp = self.request_fingerprint(request)
|
||||
# This returns the number of values added, zero if already exists.
|
||||
added = self.server.sadd(self.key, fp)
|
||||
return added == 0
|
||||
|
||||
def request_fingerprint(self, request):
|
||||
"""Returns a fingerprint for a given request.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
request : scrapy.http.Request
|
||||
|
||||
Returns
|
||||
-------
|
||||
str
|
||||
|
||||
"""
|
||||
fingerprint_data = {
|
||||
"method": to_unicode(request.method),
|
||||
"url": canonicalize_url(request.url),
|
||||
"body": (request.body or b"").hex(),
|
||||
}
|
||||
fingerprint_json = json.dumps(fingerprint_data, sort_keys=True)
|
||||
return hashlib.sha1(fingerprint_json.encode()).hexdigest()
|
||||
|
||||
@classmethod
|
||||
def from_spider(cls, spider):
|
||||
settings = spider.settings
|
||||
server = get_redis_from_settings(settings)
|
||||
dupefilter_key = settings.get(
|
||||
"SCHEDULER_DUPEFILTER_KEY", defaults.SCHEDULER_DUPEFILTER_KEY
|
||||
)
|
||||
key = dupefilter_key % {"spider": spider.name}
|
||||
debug = settings.getbool("DUPEFILTER_DEBUG")
|
||||
return cls(server, key=key, debug=debug)
|
||||
|
||||
def close(self, reason=""):
|
||||
"""Delete data on close. Called by Scrapy's scheduler.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
reason : str, optional
|
||||
|
||||
"""
|
||||
self.clear()
|
||||
|
||||
def clear(self):
|
||||
"""Clears fingerprints data."""
|
||||
self.server.delete(self.key)
|
||||
|
||||
def log(self, request, spider):
|
||||
"""Logs given request.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
request : scrapy.http.Request
|
||||
spider : scrapy.spiders.Spider
|
||||
|
||||
"""
|
||||
if self.debug:
|
||||
msg = "Filtered duplicate request: %(request)s"
|
||||
self.logger.debug(msg, {"request": request}, extra={"spider": spider})
|
||||
elif self.logdupes:
|
||||
msg = (
|
||||
"Filtered duplicate request %(request)s"
|
||||
" - no more duplicates will be shown"
|
||||
" (see DUPEFILTER_DEBUG to show all duplicates)"
|
||||
)
|
||||
self.logger.debug(msg, {"request": request}, extra={"spider": spider})
|
||||
self.logdupes = False
|
||||
@@ -0,0 +1,14 @@
|
||||
"""A pickle wrapper module with protocol=-1 by default."""
|
||||
|
||||
try:
|
||||
import cPickle as pickle # PY2
|
||||
except ImportError:
|
||||
import pickle
|
||||
|
||||
|
||||
def loads(s):
|
||||
return pickle.loads(s)
|
||||
|
||||
|
||||
def dumps(obj):
|
||||
return pickle.dumps(obj, protocol=-1)
|
||||
@@ -0,0 +1,73 @@
|
||||
from scrapy.utils.misc import load_object
|
||||
from scrapy.utils.serialize import ScrapyJSONEncoder
|
||||
from twisted.internet.threads import deferToThread
|
||||
|
||||
from . import connection, defaults
|
||||
|
||||
default_serialize = ScrapyJSONEncoder().encode
|
||||
|
||||
|
||||
class RedisPipeline:
|
||||
"""Pushes serialized item into a redis list/queue
|
||||
|
||||
Settings
|
||||
--------
|
||||
REDIS_ITEMS_KEY : str
|
||||
Redis key where to store items.
|
||||
REDIS_ITEMS_SERIALIZER : str
|
||||
Object path to serializer function.
|
||||
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self, server, key=defaults.PIPELINE_KEY, serialize_func=default_serialize
|
||||
):
|
||||
"""Initialize pipeline.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
server : StrictRedis
|
||||
Redis client instance.
|
||||
key : str
|
||||
Redis key where to store items.
|
||||
serialize_func : callable
|
||||
Items serializer function.
|
||||
|
||||
"""
|
||||
self.server = server
|
||||
self.key = key
|
||||
self.serialize = serialize_func
|
||||
|
||||
@classmethod
|
||||
def from_settings(cls, settings):
|
||||
params = {
|
||||
"server": connection.from_settings(settings),
|
||||
}
|
||||
if settings.get("REDIS_ITEMS_KEY"):
|
||||
params["key"] = settings["REDIS_ITEMS_KEY"]
|
||||
if settings.get("REDIS_ITEMS_SERIALIZER"):
|
||||
params["serialize_func"] = load_object(settings["REDIS_ITEMS_SERIALIZER"])
|
||||
|
||||
return cls(**params)
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler):
|
||||
return cls.from_settings(crawler.settings)
|
||||
|
||||
def process_item(self, item, spider):
|
||||
return deferToThread(self._process_item, item, spider)
|
||||
|
||||
def _process_item(self, item, spider):
|
||||
key = self.item_key(item, spider)
|
||||
data = self.serialize(item)
|
||||
self.server.rpush(key, data)
|
||||
return item
|
||||
|
||||
def item_key(self, item, spider):
|
||||
"""Returns redis key based on given spider.
|
||||
|
||||
Override this function to use a different key depending on the item
|
||||
and/or spider.
|
||||
|
||||
"""
|
||||
return self.key % {"spider": spider.name}
|
||||
@@ -0,0 +1,155 @@
|
||||
try:
|
||||
from scrapy.utils.request import request_from_dict
|
||||
except ImportError:
|
||||
from scrapy.utils.reqser import request_to_dict, request_from_dict
|
||||
|
||||
from . import picklecompat
|
||||
|
||||
|
||||
class Base:
|
||||
"""Per-spider base queue class"""
|
||||
|
||||
def __init__(self, server, spider, key, serializer=None):
|
||||
"""Initialize per-spider redis queue.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
server : StrictRedis
|
||||
Redis client instance.
|
||||
spider : Spider
|
||||
Scrapy spider instance.
|
||||
key: str
|
||||
Redis key where to put and get messages.
|
||||
serializer : object
|
||||
Serializer object with ``loads`` and ``dumps`` methods.
|
||||
|
||||
"""
|
||||
if serializer is None:
|
||||
# Backward compatibility.
|
||||
# TODO: deprecate pickle.
|
||||
serializer = picklecompat
|
||||
if not hasattr(serializer, "loads"):
|
||||
raise TypeError(
|
||||
f"serializer does not implement 'loads' function: {serializer}"
|
||||
)
|
||||
if not hasattr(serializer, "dumps"):
|
||||
raise TypeError(
|
||||
f"serializer does not implement 'dumps' function: {serializer}"
|
||||
)
|
||||
|
||||
self.server = server
|
||||
self.spider = spider
|
||||
self.key = key % {"spider": spider.name}
|
||||
self.serializer = serializer
|
||||
|
||||
def _encode_request(self, request):
|
||||
"""Encode a request object"""
|
||||
try:
|
||||
obj = request.to_dict(spider=self.spider)
|
||||
except AttributeError:
|
||||
obj = request_to_dict(request, self.spider)
|
||||
return self.serializer.dumps(obj)
|
||||
|
||||
def _decode_request(self, encoded_request):
|
||||
"""Decode an request previously encoded"""
|
||||
obj = self.serializer.loads(encoded_request)
|
||||
return request_from_dict(obj, spider=self.spider)
|
||||
|
||||
def __len__(self):
|
||||
"""Return the length of the queue"""
|
||||
raise NotImplementedError
|
||||
|
||||
def push(self, request):
|
||||
"""Push a request"""
|
||||
raise NotImplementedError
|
||||
|
||||
def pop(self, timeout=0):
|
||||
"""Pop a request"""
|
||||
raise NotImplementedError
|
||||
|
||||
def clear(self):
|
||||
"""Clear queue/stack"""
|
||||
self.server.delete(self.key)
|
||||
|
||||
|
||||
class FifoQueue(Base):
|
||||
"""Per-spider FIFO queue"""
|
||||
|
||||
def __len__(self):
|
||||
"""Return the length of the queue"""
|
||||
return self.server.llen(self.key)
|
||||
|
||||
def push(self, request):
|
||||
"""Push a request"""
|
||||
self.server.lpush(self.key, self._encode_request(request))
|
||||
|
||||
def pop(self, timeout=0):
|
||||
"""Pop a request"""
|
||||
if timeout > 0:
|
||||
data = self.server.brpop(self.key, timeout)
|
||||
if isinstance(data, tuple):
|
||||
data = data[1]
|
||||
else:
|
||||
data = self.server.rpop(self.key)
|
||||
if data:
|
||||
return self._decode_request(data)
|
||||
|
||||
|
||||
class PriorityQueue(Base):
|
||||
"""Per-spider priority queue abstraction using redis' sorted set"""
|
||||
|
||||
def __len__(self):
|
||||
"""Return the length of the queue"""
|
||||
return self.server.zcard(self.key)
|
||||
|
||||
def push(self, request):
|
||||
"""Push a request"""
|
||||
data = self._encode_request(request)
|
||||
score = -request.priority
|
||||
# We don't use zadd method as the order of arguments change depending on
|
||||
# whether the class is Redis or StrictRedis, and the option of using
|
||||
# kwargs only accepts strings, not bytes.
|
||||
self.server.execute_command("ZADD", self.key, score, data)
|
||||
|
||||
def pop(self, timeout=0):
|
||||
"""
|
||||
Pop a request
|
||||
timeout not support in this queue class
|
||||
"""
|
||||
# use atomic range/remove using multi/exec
|
||||
pipe = self.server.pipeline()
|
||||
pipe.multi()
|
||||
pipe.zrange(self.key, 0, 0).zremrangebyrank(self.key, 0, 0)
|
||||
results, count = pipe.execute()
|
||||
if results:
|
||||
return self._decode_request(results[0])
|
||||
|
||||
|
||||
class LifoQueue(Base):
|
||||
"""Per-spider LIFO queue."""
|
||||
|
||||
def __len__(self):
|
||||
"""Return the length of the stack"""
|
||||
return self.server.llen(self.key)
|
||||
|
||||
def push(self, request):
|
||||
"""Push a request"""
|
||||
self.server.lpush(self.key, self._encode_request(request))
|
||||
|
||||
def pop(self, timeout=0):
|
||||
"""Pop a request"""
|
||||
if timeout > 0:
|
||||
data = self.server.blpop(self.key, timeout)
|
||||
if isinstance(data, tuple):
|
||||
data = data[1]
|
||||
else:
|
||||
data = self.server.lpop(self.key)
|
||||
|
||||
if data:
|
||||
return self._decode_request(data)
|
||||
|
||||
|
||||
# TODO: Deprecate the use of these names.
|
||||
SpiderQueue = FifoQueue
|
||||
SpiderStack = LifoQueue
|
||||
SpiderPriorityQueue = PriorityQueue
|
||||
@@ -0,0 +1,182 @@
|
||||
import importlib
|
||||
|
||||
from scrapy.utils.misc import load_object
|
||||
|
||||
from . import connection, defaults
|
||||
|
||||
|
||||
# TODO: add SCRAPY_JOB support.
|
||||
class Scheduler:
|
||||
"""Redis-based scheduler
|
||||
|
||||
Settings
|
||||
--------
|
||||
SCHEDULER_PERSIST : bool (default: False)
|
||||
Whether to persist or clear redis queue.
|
||||
SCHEDULER_FLUSH_ON_START : bool (default: False)
|
||||
Whether to flush redis queue on start.
|
||||
SCHEDULER_IDLE_BEFORE_CLOSE : int (default: 0)
|
||||
How many seconds to wait before closing if no message is received.
|
||||
SCHEDULER_QUEUE_KEY : str
|
||||
Scheduler redis key.
|
||||
SCHEDULER_QUEUE_CLASS : str
|
||||
Scheduler queue class.
|
||||
SCHEDULER_DUPEFILTER_KEY : str
|
||||
Scheduler dupefilter redis key.
|
||||
SCHEDULER_DUPEFILTER_CLASS : str
|
||||
Scheduler dupefilter class.
|
||||
SCHEDULER_SERIALIZER : str
|
||||
Scheduler serializer.
|
||||
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
server,
|
||||
persist=False,
|
||||
flush_on_start=False,
|
||||
queue_key=defaults.SCHEDULER_QUEUE_KEY,
|
||||
queue_cls=defaults.SCHEDULER_QUEUE_CLASS,
|
||||
dupefilter=None,
|
||||
dupefilter_key=defaults.SCHEDULER_DUPEFILTER_KEY,
|
||||
dupefilter_cls=defaults.SCHEDULER_DUPEFILTER_CLASS,
|
||||
idle_before_close=0,
|
||||
serializer=None,
|
||||
):
|
||||
"""Initialize scheduler.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
server : Redis
|
||||
The redis server instance.
|
||||
persist : bool
|
||||
Whether to flush requests when closing. Default is False.
|
||||
flush_on_start : bool
|
||||
Whether to flush requests on start. Default is False.
|
||||
queue_key : str
|
||||
Requests queue key.
|
||||
queue_cls : str
|
||||
Importable path to the queue class.
|
||||
dupefilter: Dupefilter
|
||||
Custom dupefilter instance.
|
||||
dupefilter_key : str
|
||||
Duplicates filter key.
|
||||
dupefilter_cls : str
|
||||
Importable path to the dupefilter class.
|
||||
idle_before_close : int
|
||||
Timeout before giving up.
|
||||
|
||||
"""
|
||||
if idle_before_close < 0:
|
||||
raise TypeError("idle_before_close cannot be negative")
|
||||
|
||||
self.server = server
|
||||
self.persist = persist
|
||||
self.flush_on_start = flush_on_start
|
||||
self.queue_key = queue_key
|
||||
self.queue_cls = queue_cls
|
||||
self.df = dupefilter
|
||||
self.dupefilter_cls = dupefilter_cls
|
||||
self.dupefilter_key = dupefilter_key
|
||||
self.idle_before_close = idle_before_close
|
||||
self.serializer = serializer
|
||||
self.stats = None
|
||||
|
||||
def __len__(self):
|
||||
return len(self.queue)
|
||||
|
||||
@classmethod
|
||||
def from_settings(cls, settings):
|
||||
kwargs = {
|
||||
"persist": settings.getbool("SCHEDULER_PERSIST"),
|
||||
"flush_on_start": settings.getbool("SCHEDULER_FLUSH_ON_START"),
|
||||
"idle_before_close": settings.getint("SCHEDULER_IDLE_BEFORE_CLOSE"),
|
||||
}
|
||||
|
||||
# If these values are missing, it means we want to use the defaults.
|
||||
optional = {
|
||||
# TODO: Use custom prefixes for this settings to note that are
|
||||
# specific to scrapy-redis.
|
||||
"queue_key": "SCHEDULER_QUEUE_KEY",
|
||||
"queue_cls": "SCHEDULER_QUEUE_CLASS",
|
||||
"dupefilter_key": "SCHEDULER_DUPEFILTER_KEY",
|
||||
# We use the default setting name to keep compatibility.
|
||||
"dupefilter_cls": "DUPEFILTER_CLASS",
|
||||
"serializer": "SCHEDULER_SERIALIZER",
|
||||
}
|
||||
for name, setting_name in optional.items():
|
||||
val = settings.get(setting_name)
|
||||
if val:
|
||||
kwargs[name] = val
|
||||
|
||||
dupefilter_cls = load_object(kwargs["dupefilter_cls"])
|
||||
if not hasattr(dupefilter_cls, "from_spider"):
|
||||
kwargs["dupefilter"] = dupefilter_cls.from_settings(settings)
|
||||
|
||||
# Support serializer as a path to a module.
|
||||
if isinstance(kwargs.get("serializer"), str):
|
||||
kwargs["serializer"] = importlib.import_module(kwargs["serializer"])
|
||||
|
||||
server = connection.from_settings(settings)
|
||||
# Ensure the connection is working.
|
||||
server.ping()
|
||||
|
||||
return cls(server=server, **kwargs)
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler):
|
||||
instance = cls.from_settings(crawler.settings)
|
||||
# FIXME: for now, stats are only supported from this constructor
|
||||
instance.stats = crawler.stats
|
||||
return instance
|
||||
|
||||
def open(self, spider):
|
||||
self.spider = spider
|
||||
|
||||
try:
|
||||
self.queue = load_object(self.queue_cls)(
|
||||
server=self.server,
|
||||
spider=spider,
|
||||
key=self.queue_key % {"spider": spider.name},
|
||||
serializer=self.serializer,
|
||||
)
|
||||
except TypeError as e:
|
||||
raise ValueError(
|
||||
f"Failed to instantiate queue class '{self.queue_cls}': {e}"
|
||||
)
|
||||
|
||||
if not self.df:
|
||||
self.df = load_object(self.dupefilter_cls).from_spider(spider)
|
||||
|
||||
if self.flush_on_start:
|
||||
self.flush()
|
||||
# notice if there are requests already in the queue to resume the crawl
|
||||
if len(self.queue):
|
||||
spider.log(f"Resuming crawl ({len(self.queue)} requests scheduled)")
|
||||
|
||||
def close(self, reason):
|
||||
if not self.persist:
|
||||
self.flush()
|
||||
|
||||
def flush(self):
|
||||
self.df.clear()
|
||||
self.queue.clear()
|
||||
|
||||
def enqueue_request(self, request):
|
||||
if not request.dont_filter and self.df.request_seen(request):
|
||||
self.df.log(request, self.spider)
|
||||
return False
|
||||
if self.stats:
|
||||
self.stats.inc_value("scheduler/enqueued/redis", spider=self.spider)
|
||||
self.queue.push(request)
|
||||
return True
|
||||
|
||||
def next_request(self):
|
||||
block_pop_timeout = self.idle_before_close
|
||||
request = self.queue.pop(block_pop_timeout)
|
||||
if request and self.stats:
|
||||
self.stats.inc_value("scheduler/dequeued/redis", spider=self.spider)
|
||||
return request
|
||||
|
||||
def has_pending_requests(self):
|
||||
return len(self) > 0
|
||||
@@ -0,0 +1,297 @@
|
||||
import json
|
||||
import time
|
||||
from collections.abc import Iterable
|
||||
|
||||
from scrapy import FormRequest, signals
|
||||
from scrapy import version_info as scrapy_version
|
||||
from scrapy.exceptions import DontCloseSpider
|
||||
from scrapy.spiders import CrawlSpider, Spider
|
||||
|
||||
from scrapy_redis.utils import TextColor
|
||||
|
||||
from . import connection, defaults
|
||||
from .utils import bytes_to_str, is_dict
|
||||
|
||||
|
||||
class RedisMixin:
|
||||
"""Mixin class to implement reading urls from a redis queue."""
|
||||
|
||||
redis_key = None
|
||||
redis_batch_size = None
|
||||
redis_encoding = None
|
||||
|
||||
# Redis client placeholder.
|
||||
server = None
|
||||
|
||||
# Idle start time
|
||||
spider_idle_start_time = int(time.time())
|
||||
max_idle_time = None
|
||||
|
||||
def start_requests(self):
|
||||
"""Returns a batch of start requests from redis."""
|
||||
return self.next_requests()
|
||||
|
||||
def setup_redis(self, crawler=None):
|
||||
"""Setup redis connection and idle signal.
|
||||
|
||||
This should be called after the spider has set its crawler object.
|
||||
"""
|
||||
if self.server is not None:
|
||||
return
|
||||
|
||||
if crawler is None:
|
||||
# We allow optional crawler argument to keep backwards
|
||||
# compatibility.
|
||||
# XXX: Raise a deprecation warning.
|
||||
crawler = getattr(self, "crawler", None)
|
||||
|
||||
if crawler is None:
|
||||
raise ValueError("crawler is required")
|
||||
|
||||
settings = crawler.settings
|
||||
|
||||
if self.redis_key is None:
|
||||
self.redis_key = settings.get(
|
||||
"REDIS_START_URLS_KEY",
|
||||
defaults.START_URLS_KEY,
|
||||
)
|
||||
|
||||
self.redis_key = self.redis_key % {"name": self.name}
|
||||
|
||||
if not self.redis_key.strip():
|
||||
raise ValueError("redis_key must not be empty")
|
||||
|
||||
if self.redis_batch_size is None:
|
||||
self.redis_batch_size = settings.getint(
|
||||
"CONCURRENT_REQUESTS", defaults.REDIS_CONCURRENT_REQUESTS
|
||||
)
|
||||
|
||||
try:
|
||||
self.redis_batch_size = int(self.redis_batch_size)
|
||||
except (TypeError, ValueError):
|
||||
raise ValueError("redis_batch_size must be an integer")
|
||||
|
||||
if self.redis_encoding is None:
|
||||
self.redis_encoding = settings.get(
|
||||
"REDIS_ENCODING", defaults.REDIS_ENCODING
|
||||
)
|
||||
|
||||
self.logger.info(
|
||||
"Reading start URLs from redis key '%(redis_key)s' "
|
||||
"(batch size: %(redis_batch_size)s, encoding: %(redis_encoding)s)",
|
||||
self.__dict__,
|
||||
)
|
||||
|
||||
self.server = connection.from_settings(crawler.settings)
|
||||
|
||||
if settings.getbool("REDIS_START_URLS_AS_SET", defaults.START_URLS_AS_SET):
|
||||
self.fetch_data = self.server.spop
|
||||
self.count_size = self.server.scard
|
||||
elif settings.getbool("REDIS_START_URLS_AS_ZSET", defaults.START_URLS_AS_ZSET):
|
||||
self.fetch_data = self.pop_priority_queue
|
||||
self.count_size = self.server.zcard
|
||||
else:
|
||||
self.fetch_data = self.pop_list_queue
|
||||
self.count_size = self.server.llen
|
||||
|
||||
if self.max_idle_time is None:
|
||||
self.max_idle_time = settings.get(
|
||||
"MAX_IDLE_TIME_BEFORE_CLOSE", defaults.MAX_IDLE_TIME
|
||||
)
|
||||
|
||||
try:
|
||||
self.max_idle_time = int(self.max_idle_time)
|
||||
except (TypeError, ValueError):
|
||||
raise ValueError("max_idle_time must be an integer")
|
||||
|
||||
# The idle signal is called when the spider has no requests left,
|
||||
# that's when we will schedule new requests from redis queue
|
||||
crawler.signals.connect(self.spider_idle, signal=signals.spider_idle)
|
||||
|
||||
def pop_list_queue(self, redis_key, batch_size):
|
||||
with self.server.pipeline() as pipe:
|
||||
pipe.lrange(redis_key, 0, batch_size - 1)
|
||||
pipe.ltrim(redis_key, batch_size, -1)
|
||||
datas, _ = pipe.execute()
|
||||
return datas
|
||||
|
||||
def pop_priority_queue(self, redis_key, batch_size):
|
||||
with self.server.pipeline() as pipe:
|
||||
pipe.zrevrange(redis_key, 0, batch_size - 1)
|
||||
pipe.zremrangebyrank(redis_key, -batch_size, -1)
|
||||
datas, _ = pipe.execute()
|
||||
return datas
|
||||
|
||||
def next_requests(self):
|
||||
"""Returns a request to be scheduled or none."""
|
||||
# XXX: Do we need to use a timeout here?
|
||||
found = 0
|
||||
datas = self.fetch_data(self.redis_key, self.redis_batch_size)
|
||||
for data in datas:
|
||||
reqs = self.make_request_from_data(data)
|
||||
if isinstance(reqs, Iterable):
|
||||
for req in reqs:
|
||||
yield req
|
||||
# XXX: should be here?
|
||||
found += 1
|
||||
self.logger.info(f"start req url:{req.url}")
|
||||
elif reqs:
|
||||
yield reqs
|
||||
found += 1
|
||||
else:
|
||||
self.logger.debug(f"Request not made from data: {data}")
|
||||
|
||||
if found:
|
||||
self.logger.debug(f"Read {found} requests from '{self.redis_key}'")
|
||||
|
||||
def make_request_from_data(self, data):
|
||||
"""Returns a `Request` instance for data coming from Redis.
|
||||
|
||||
Overriding this function to support the `json` requested `data` that contains
|
||||
`url` ,`meta` and other optional parameters. `meta` is a nested json which contains sub-data.
|
||||
|
||||
Along with:
|
||||
After accessing the data, sending the FormRequest with `url`, `meta` and addition `formdata`, `method`
|
||||
|
||||
For example:
|
||||
|
||||
.. code:: json
|
||||
|
||||
{
|
||||
"url": "https://example.com",
|
||||
"meta": {
|
||||
"job-id":"123xsd",
|
||||
"start-date":"dd/mm/yy",
|
||||
},
|
||||
"url_cookie_key":"fertxsas",
|
||||
"method":"POST",
|
||||
}
|
||||
|
||||
If `url` is empty, return `[]`. So you should verify the `url` in the data.
|
||||
If `method` is empty, the request object will set method to 'GET', optional.
|
||||
If `meta` is empty, the request object will set `meta` to an empty dictionary, optional.
|
||||
|
||||
This json supported data can be accessed from 'scrapy.spider' through response.
|
||||
'request.url', 'request.meta', 'request.cookies', 'request.method'
|
||||
|
||||
Parameters
|
||||
----------
|
||||
data : bytes
|
||||
Message from redis.
|
||||
|
||||
"""
|
||||
formatted_data = bytes_to_str(data, self.redis_encoding)
|
||||
|
||||
if is_dict(formatted_data):
|
||||
parameter = json.loads(formatted_data)
|
||||
else:
|
||||
self.logger.warning(
|
||||
f"{TextColor.WARNING}WARNING: String request is deprecated, please use JSON data format. "
|
||||
f"Detail information, please check https://github.com/rmax/scrapy-redis#features{TextColor.ENDC}"
|
||||
)
|
||||
return FormRequest(formatted_data, dont_filter=True)
|
||||
|
||||
if parameter.get("url", None) is None:
|
||||
self.logger.warning(
|
||||
f"{TextColor.WARNING}The data from Redis has no url key in push data{TextColor.ENDC}"
|
||||
)
|
||||
return []
|
||||
|
||||
url = parameter.pop("url")
|
||||
method = parameter.pop("method").upper() if "method" in parameter else "GET"
|
||||
metadata = parameter.pop("meta") if "meta" in parameter else {}
|
||||
|
||||
return FormRequest(
|
||||
url, dont_filter=True, method=method, formdata=parameter, meta=metadata
|
||||
)
|
||||
|
||||
def schedule_next_requests(self):
|
||||
"""Schedules a request if available"""
|
||||
# TODO: While there is capacity, schedule a batch of redis requests.
|
||||
for req in self.next_requests():
|
||||
# see https://github.com/scrapy/scrapy/issues/5994
|
||||
if scrapy_version >= (2, 6):
|
||||
self.crawler.engine.crawl(req)
|
||||
else:
|
||||
self.crawler.engine.crawl(req, spider=self)
|
||||
|
||||
def spider_idle(self):
|
||||
"""
|
||||
Schedules a request if available, otherwise waits.
|
||||
or close spider when waiting seconds > MAX_IDLE_TIME_BEFORE_CLOSE.
|
||||
MAX_IDLE_TIME_BEFORE_CLOSE will not affect SCHEDULER_IDLE_BEFORE_CLOSE.
|
||||
"""
|
||||
if self.server is not None and self.count_size(self.redis_key) > 0:
|
||||
self.spider_idle_start_time = int(time.time())
|
||||
|
||||
self.schedule_next_requests()
|
||||
|
||||
idle_time = int(time.time()) - self.spider_idle_start_time
|
||||
if self.max_idle_time != 0 and idle_time >= self.max_idle_time:
|
||||
return
|
||||
raise DontCloseSpider
|
||||
|
||||
|
||||
class RedisSpider(RedisMixin, Spider):
|
||||
"""Spider that reads urls from redis queue when idle.
|
||||
|
||||
Attributes
|
||||
----------
|
||||
redis_key : str (default: REDIS_START_URLS_KEY)
|
||||
Redis key where to fetch start URLs from..
|
||||
redis_batch_size : int (default: CONCURRENT_REQUESTS)
|
||||
Number of messages to fetch from redis on each attempt.
|
||||
redis_encoding : str (default: REDIS_ENCODING)
|
||||
Encoding to use when decoding messages from redis queue.
|
||||
|
||||
Settings
|
||||
--------
|
||||
REDIS_START_URLS_KEY : str (default: "<spider.name>:start_urls")
|
||||
Default Redis key where to fetch start URLs from..
|
||||
REDIS_START_URLS_BATCH_SIZE : int (deprecated by CONCURRENT_REQUESTS)
|
||||
Default number of messages to fetch from redis on each attempt.
|
||||
REDIS_START_URLS_AS_SET : bool (default: False)
|
||||
Use SET operations to retrieve messages from the redis queue. If False,
|
||||
the messages are retrieve using the LPOP command.
|
||||
REDIS_ENCODING : str (default: "utf-8")
|
||||
Default encoding to use when decoding messages from redis queue.
|
||||
|
||||
"""
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler, *args, **kwargs):
|
||||
obj = super().from_crawler(crawler, *args, **kwargs)
|
||||
obj.setup_redis(crawler)
|
||||
return obj
|
||||
|
||||
|
||||
class RedisCrawlSpider(RedisMixin, CrawlSpider):
|
||||
"""Spider that reads urls from redis queue when idle.
|
||||
|
||||
Attributes
|
||||
----------
|
||||
redis_key : str (default: REDIS_START_URLS_KEY)
|
||||
Redis key where to fetch start URLs from..
|
||||
redis_batch_size : int (default: CONCURRENT_REQUESTS)
|
||||
Number of messages to fetch from redis on each attempt.
|
||||
redis_encoding : str (default: REDIS_ENCODING)
|
||||
Encoding to use when decoding messages from redis queue.
|
||||
|
||||
Settings
|
||||
--------
|
||||
REDIS_START_URLS_KEY : str (default: "<spider.name>:start_urls")
|
||||
Default Redis key where to fetch start URLs from..
|
||||
REDIS_START_URLS_BATCH_SIZE : int (deprecated by CONCURRENT_REQUESTS)
|
||||
Default number of messages to fetch from redis on each attempt.
|
||||
REDIS_START_URLS_AS_SET : bool (default: True)
|
||||
Use SET operations to retrieve messages from the redis queue.
|
||||
REDIS_ENCODING : str (default: "utf-8")
|
||||
Default encoding to use when decoding messages from redis queue.
|
||||
|
||||
"""
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler, *args, **kwargs):
|
||||
obj = super().from_crawler(crawler, *args, **kwargs)
|
||||
obj.setup_redis(crawler)
|
||||
return obj
|
||||
@@ -0,0 +1,90 @@
|
||||
from datetime import datetime
|
||||
|
||||
from scrapy.statscollectors import StatsCollector
|
||||
|
||||
from .connection import from_settings as redis_from_settings
|
||||
from .defaults import SCHEDULER_PERSIST, STATS_KEY
|
||||
from .utils import convert_bytes_to_str
|
||||
|
||||
|
||||
class RedisStatsCollector(StatsCollector):
|
||||
"""
|
||||
Stats Collector based on Redis
|
||||
"""
|
||||
|
||||
def __init__(self, crawler, spider=None):
|
||||
super().__init__(crawler)
|
||||
self.server = redis_from_settings(crawler.settings)
|
||||
self.spider = spider
|
||||
self.spider_name = spider.name if spider else crawler.spidercls.name
|
||||
self.stats_key = crawler.settings.get("STATS_KEY", STATS_KEY)
|
||||
self.persist = crawler.settings.get("SCHEDULER_PERSIST", SCHEDULER_PERSIST)
|
||||
|
||||
def _get_key(self, spider=None):
|
||||
"""Return the hash name of stats"""
|
||||
if spider:
|
||||
return self.stats_key % {"spider": spider.name}
|
||||
if self.spider:
|
||||
return self.stats_key % {"spider": self.spider.name}
|
||||
return self.stats_key % {"spider": self.spider_name or "scrapy"}
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler):
|
||||
return cls(crawler)
|
||||
|
||||
@classmethod
|
||||
def from_spider(cls, spider):
|
||||
return cls(spider.crawler)
|
||||
|
||||
def get_value(self, key, default=None, spider=None):
|
||||
"""Return the value of hash stats"""
|
||||
if self.server.hexists(self._get_key(spider), key):
|
||||
return int(self.server.hget(self._get_key(spider), key))
|
||||
else:
|
||||
return default
|
||||
|
||||
def get_stats(self, spider=None):
|
||||
"""Return the all of the values of hash stats"""
|
||||
stats = self.server.hgetall(self._get_key(spider))
|
||||
if stats:
|
||||
return convert_bytes_to_str(stats)
|
||||
return {}
|
||||
|
||||
def set_value(self, key, value, spider=None):
|
||||
"""Set the value according to hash key of stats"""
|
||||
if isinstance(value, datetime):
|
||||
value = value.timestamp()
|
||||
self.server.hset(self._get_key(spider), key, value)
|
||||
|
||||
def set_stats(self, stats, spider=None):
|
||||
"""Set all the hash stats"""
|
||||
self.server.hmset(self._get_key(spider), stats)
|
||||
|
||||
def inc_value(self, key, count=1, start=0, spider=None):
|
||||
"""Set increment of value according to key"""
|
||||
if not self.server.hexists(self._get_key(spider), key):
|
||||
self.set_value(key, start)
|
||||
self.server.hincrby(self._get_key(spider), key, count)
|
||||
|
||||
def max_value(self, key, value, spider=None):
|
||||
"""Set max value between current and new value"""
|
||||
self.set_value(key, max(self.get_value(key, value), value))
|
||||
|
||||
def min_value(self, key, value, spider=None):
|
||||
"""Set min value between current and new value"""
|
||||
self.set_value(key, min(self.get_value(key, value), value))
|
||||
|
||||
def clear_stats(self, spider=None):
|
||||
"""Clear all the hash stats"""
|
||||
self.server.delete(self._get_key(spider))
|
||||
|
||||
def open_spider(self, spider):
|
||||
"""Set spider to self"""
|
||||
if spider:
|
||||
self.spider = spider
|
||||
|
||||
def close_spider(self, spider, reason):
|
||||
"""Clear spider and clear stats"""
|
||||
self.spider = None
|
||||
if not self.persist:
|
||||
self.clear_stats(spider)
|
||||
@@ -0,0 +1,44 @@
|
||||
import json
|
||||
from json import JSONDecodeError
|
||||
|
||||
import six
|
||||
|
||||
|
||||
class TextColor:
|
||||
HEADER = "\033[95m"
|
||||
OKBLUE = "\033[94m"
|
||||
OKCYAN = "\033[96m"
|
||||
OKGREEN = "\033[92m"
|
||||
WARNING = "\033[93m"
|
||||
FAIL = "\033[91m"
|
||||
ENDC = "\033[0m"
|
||||
BOLD = "\033[1m"
|
||||
UNDERLINE = "\033[4m"
|
||||
|
||||
|
||||
def bytes_to_str(s, encoding="utf-8"):
|
||||
"""Returns a str if a bytes object is given."""
|
||||
if six.PY3 and isinstance(s, bytes):
|
||||
return s.decode(encoding)
|
||||
return s
|
||||
|
||||
|
||||
def is_dict(string_content):
|
||||
"""Try load string_content as json, if failed, return False, else return True."""
|
||||
try:
|
||||
json.loads(string_content)
|
||||
except JSONDecodeError:
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def convert_bytes_to_str(data, encoding="utf-8"):
|
||||
"""Convert a dict's keys & values from `bytes` to `str`
|
||||
or convert bytes to str"""
|
||||
if isinstance(data, bytes):
|
||||
return data.decode(encoding)
|
||||
if isinstance(data, dict):
|
||||
return dict(map(convert_bytes_to_str, data.items()))
|
||||
elif isinstance(data, tuple):
|
||||
return map(convert_bytes_to_str, data)
|
||||
return data
|
||||
Reference in New Issue
Block a user