Add dramatiq-abort support and extend scrapy CloseSpider to abort the dramatiq worker on stale spiders

Integrate `dramatiq-abort` with Redis backend for task abortion and enhance the `CloseSpider` extension to abort jobs with no recent activity. These changes improve task management and error handling for idle spiders in the scraping process. Updated dependencies and project settings accordingly.
This commit is contained in:
mhdzumair
2024-12-15 23:01:24 +05:30
parent 402c86b907
commit 863247018d
5 changed files with 58 additions and 22 deletions
+1
View File
@@ -48,6 +48,7 @@ aioseedrcc = "*"
pytz = "*"
aiohttp-socks = "*"
socksio = "*"
dramatiq-abort = "*"
[dev-packages]
pysocks = "*"
Generated
+27 -19
View File
@@ -1,7 +1,7 @@
{
"_meta": {
"hash": {
"sha256": "4c1f290f209eb78b77a8fb350c4e502aba3bcd06355e3a7291152849af96b459"
"sha256": "163b3b6a3d97b3f660a2a203ee1fc720df6d8946ef07e0c8309965ba05f7963e"
},
"pipfile-spec": 6,
"requires": {
@@ -146,11 +146,11 @@
},
"aiosignal": {
"hashes": [
"sha256:54cd96e15e1649b75d6c87526a6ff0b6c1b0dd3459f43d9ca11d48c339b68cfc",
"sha256:f8376fb07dd1e86a584e4fcdec80b36b7f81aac666ebc724e2c090300dd83b17"
"sha256:45cde58e409a301715980c2b01d0c28bdde3770d8290b5eb2173759d9acb31a5",
"sha256:a8c255c66fafb1e499c9351d0bf32ff2d8a0321595ebac3b93713656d2436f54"
],
"markers": "python_version >= '3.7'",
"version": "==1.3.1"
"markers": "python_version >= '3.9'",
"version": "==1.3.2"
},
"aiowebdav": {
"hashes": [
@@ -237,11 +237,11 @@
},
"certifi": {
"hashes": [
"sha256:922820b53db7a7257ffbda3f597266d435245903d80737e34f8a45ff3e3230d8",
"sha256:bec941d2aa8195e248a60b31ff9f0558284cf01a52591ceda73ea9afffd69fd9"
"sha256:1275f7a45be9464efc1173084eaa30f866fe2e47d389406136d332ed4967ec56",
"sha256:b650d30f370c2b724812bee08008be0c4163b163ddaec3f2546c1caf65f191db"
],
"markers": "python_version >= '3.6'",
"version": "==2024.8.30"
"version": "==2024.12.14"
},
"cffi": {
"hashes": [
@@ -552,6 +552,14 @@
"markers": "python_version >= '3.8'",
"version": "==1.17.1"
},
"dramatiq-abort": {
"hashes": [
"sha256:795ee9a37d884d70e63b758031017cd1f7c33f40db5f4901aaa95fcb659c13c0",
"sha256:d025e358875b437346cc2f665525735b1a079ce7dcac9c3dc70853fa798efd69"
],
"index": "pypi",
"version": "==1.2.1"
},
"fastapi": {
"hashes": [
"sha256:9ec46f7addc14ea472958a96aae5b5de65f39721a46aaf5705c480d9a8b76654",
@@ -995,7 +1003,7 @@
"humanize": {
"git": "git+https://github.com/python-humanize/humanize.git",
"markers": "python_version >= '3.9'",
"ref": "767c7cb0b21f27ade9ebfc84d24b5eecc07acd43"
"ref": "db482bbbc2f52f86a804d32acfccf0264ebdf5c6"
},
"hyperlink": {
"hashes": [
@@ -1899,12 +1907,12 @@
},
"pydantic-settings": {
"hashes": [
"sha256:7fb0637c786a558d3103436278a7c4f1cfd29ba8973238a50c5bb9a55387da87",
"sha256:e0f92546d8a9923cb8941689abf85d6601a8c19a23e97a34b2964a2e3f813ca0"
"sha256:ac4bfd4a36831a48dbf8b2d9325425b549a0a6f18cea118436d728eb4f1c4d66",
"sha256:e00c05d5fa6cbbb227c84bd7487c5c1065084119b750df7c8c1a554aed236eb5"
],
"index": "pypi",
"markers": "python_version >= '3.8'",
"version": "==2.6.1"
"version": "==2.7.0"
},
"pydispatcher": {
"hashes": [
@@ -2641,11 +2649,11 @@
"standard"
],
"hashes": [
"sha256:82ad92fd58da0d12af7482ecdb5f2470a04c9c9a53ced65b9bbb4a205377602e",
"sha256:ee9519c246a72b1c084cea8d3b44ed6026e78a4a309cbedae9c37e4cb9fbb175"
"sha256:023dc038422502fa28a09c7a30bf2b6991512da7dcdb8fd35fe57cfc154126f4",
"sha256:404051050cd7e905de2c9a7e61790943440b3416f49cb409f965d9dcd0fa73e9"
],
"markers": "python_version >= '3.8'",
"version": "==0.32.1"
"markers": "python_version >= '3.9'",
"version": "==0.34.0"
},
"uvloop": {
"hashes": [
@@ -3049,11 +3057,11 @@
},
"certifi": {
"hashes": [
"sha256:922820b53db7a7257ffbda3f597266d435245903d80737e34f8a45ff3e3230d8",
"sha256:bec941d2aa8195e248a60b31ff9f0558284cf01a52591ceda73ea9afffd69fd9"
"sha256:1275f7a45be9464efc1173084eaa30f866fe2e47d389406136d332ed4967ec56",
"sha256:b650d30f370c2b724812bee08008be0c4163b163ddaec3f2546c1caf65f191db"
],
"markers": "python_version >= '3.6'",
"version": "==2024.8.30"
"version": "==2024.12.14"
},
"charset-normalizer": {
"hashes": [
+5
View File
@@ -9,7 +9,10 @@ from dramatiq.middleware import (
Callbacks,
Pipelines,
Prometheus,
CurrentMessage,
)
from dramatiq_abort import Abortable
from dramatiq_abort.backends import RedisBackend
from api.middleware import MaxTasksPerChild, Retries, TaskManager
from db.config import settings
@@ -27,5 +30,7 @@ redis_broker.middleware = [
AsyncIO(),
MaxTasksPerChild(settings.worker_max_tasks_per_child),
TaskManager(),
CurrentMessage(),
Abortable(backend=RedisBackend.from_url(settings.redis_url)),
]
dramatiq.set_broker(redis_broker)
+24
View File
@@ -0,0 +1,24 @@
from dramatiq.middleware import CurrentMessage
from dramatiq_abort import abort
from scrapy import Spider
from scrapy.extensions.closespider import CloseSpider
class CloseSpiderExtended(CloseSpider):
def _count_items_produced(self, spider: Spider) -> None:
if self.items_in_period >= 1:
self.items_in_period = 0
else:
spider.logger.info(
f"Closing spider since no items were produced in the last "
f"{self.timeout_no_item} seconds."
)
dramatiq_message = CurrentMessage.get_current_message()
if dramatiq_message:
spider.logger.info(
f"Aborting message {dramatiq_message.message_id} due to no items produced."
)
abort(dramatiq_message.message_id, abort_ttl=0)
assert self.crawler.engine
self.crawler.engine.close_spider(spider, "closespider_timeout_no_item")
+1 -3
View File
@@ -59,7 +59,7 @@ DOWNLOADER_MIDDLEWARES = {
}
EXTENSIONS = {
"scrapy.extensions.closespider.CloseSpider": 100,
"mediafusion_scrapy.extensions.CloseSpiderExtended": 100,
}
# Configure item pipelines
@@ -92,8 +92,6 @@ HTTPCACHE_IGNORE_HTTP_CODES = [500, 502, 503, 504, 522, 524, 408, 429]
HTTPCACHE_STORAGE = "scrapy.extensions.httpcache.FilesystemCacheStorage"
# HTTPCACHE_POLICY = "scrapy.extensions.httpcache.RFC2616Policy"
# Set settings whose default value is deprecated to a future-proof value
REQUEST_FINGERPRINTER_IMPLEMENTATION = "2.12"
TWISTED_REACTOR = "twisted.internet.asyncioreactor.AsyncioSelectorReactor"
FEED_EXPORT_ENCODING = "utf-8"
LOG_LEVEL = "INFO"