Skip to content

Instantly share code, notes, and snippets.

@abevoelker
Last active March 17, 2024 05:59
Show Gist options
  • Save abevoelker/10606489 to your computer and use it in GitHub Desktop.
Save abevoelker/10606489 to your computer and use it in GitHub Desktop.
scrapy RabbitMQ pipeline
# project_name/pipelines.py
from scrapy import signals
from scrapy.utils.serialize import ScrapyJSONEncoder
from scrapy.xlib.pydispatch import dispatcher
from carrot.connection import BrokerConnection
from carrot.messaging import Publisher
from twisted.internet.threads import deferToThread
import json as simplejson
import settings
class MessageQueuePipeline(object):
"""Emit processed items to a RabbitMQ exchange/queue"""
def __init__(self, host_name, port, userid, password, virtual_host, encoder_class):
self.q_connection = BrokerConnection(hostname=host_name, port=port,
userid=userid, password=password,
virtual_host=virtual_host)
self.encoder = encoder_class()
dispatcher.connect(self.spider_opened, signals.spider_opened)
dispatcher.connect(self.spider_closed, signals.spider_closed)
@classmethod
def from_settings(cls, settings):
host_name = settings.get('BROKER_HOST')
port = settings.get('BROKER_PORT')
userid = settings.get('BROKER_USERID')
password = settings.get('BROKER_PASSWORD')
virtual_host = settings.get('BROKER_VIRTUAL_HOST')
encoder_class = settings.get('MESSAGE_Q_SERIALIZER', ScrapyJSONEncoder)
return cls(host_name, port, userid, password, virtual_host, encoder_class)
def spider_opened(self, spider):
self.publisher = Publisher(connection=self.q_connection,
exchange="",
routing_key=spider.name)
def spider_closed(self, spider):
self.publisher.close()
def process_item(self, item, spider):
return deferToThread(self._process_item, item, spider)
def _process_item(self, item, spider):
self.publisher.send(self.encoder.encode(dict(item)))
return item
# project_name/settings.py
import os
BOT_NAME = 'project_name'
SPIDER_MODULES = ['project_name.spiders']
NEWSPIDER_MODULE = 'project_name.spiders'
ITEM_PIPELINES = [
'project_name.pipelines.MessageQueuePipeline',
]
try:
BROKER_HOST = os.environ['BROKER_HOST']
except KeyError:
BROKER_HOST = 'localhost'
try:
BROKER_PORT = os.environ['BROKER_PORT']
except KeyError:
BROKER_PORT = 5672
try:
BROKER_USERID = os.environ['BROKER_USERID']
except KeyError:
BROKER_USERID = 'guest'
try:
BROKER_PASSWORD = os.environ['BROKER_PASSWORD']
except KeyError:
BROKER_PASSWORD = 'guest'
try:
BROKER_VIRTUAL_HOST = os.environ['BROKER_VIRTUAL_HOST']
except KeyError:
BROKER_VIRTUAL_HOST = '/'
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment