diff --git a/Pipfile b/Pipfile index 22541fa..549f9a9 100644 --- a/Pipfile +++ b/Pipfile @@ -48,6 +48,7 @@ aioseedrcc = "*" pytz = "*" aiohttp-socks = "*" socksio = "*" +dramatiq-abort = "*" [dev-packages] pysocks = "*" diff --git a/Pipfile.lock b/Pipfile.lock index 92f99f1..e42a821 100644 --- a/Pipfile.lock +++ b/Pipfile.lock @@ -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": [ diff --git a/api/__init__.py b/api/__init__.py index dd1f770..05525ca 100644 --- a/api/__init__.py +++ b/api/__init__.py @@ -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) diff --git a/mediafusion_scrapy/extensions.py b/mediafusion_scrapy/extensions.py new file mode 100644 index 0000000..6de8572 --- /dev/null +++ b/mediafusion_scrapy/extensions.py @@ -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") diff --git a/mediafusion_scrapy/settings.py b/mediafusion_scrapy/settings.py index 7e92a15..95c353f 100644 --- a/mediafusion_scrapy/settings.py +++ b/mediafusion_scrapy/settings.py @@ -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"