Skip to content

About

A web crawling framework.

Resources

Stars

3 stars

Watchers

0 watching

Forks

Latest commit

 

History

39 Commits

Folders and files

Repository files navigation

Rubberneck

Rubberneck is an easy-to-use, multithreaded crawler framework.

Package Layout

rubberneck/
├── engine/       Runtime, routing
├── model/        Request, Response, Failure
├── scheduler/    Request queue
├── downloader/   Request -> Response
├── spider/       Response -> Item
├── pipeline/     Item consuming
├── logger/       Runtime logging
└── registry/     Component registries

Execution Flow

Rubberneck has six core runtime components: Engine, Scheduler, Downloader, Spider, Pipeline, and Logger.

Scheduler, Downloader, Spider, and Pipeline can run work in parallel; the Engine coordinates them and routes returned values between components by type.

flowchart TD
    subgraph SchedulerModule["Scheduler"]
        Enqueue["Scheduler.enqueue()"]
        Dequeue["Scheduler.dequeue()"]
        Done["Scheduler.mark_done()"]
        Failed["Scheduler.mark_failed()"]
    end
    
    subgraph EngineModule["Engine"]
        Seed["Engine.seed()"]

        subgraph EngineDownloader[" "]
            SubmitD["Engine._submit_downloader()"]
            RouteD["Engine._route_downloader_value()"]
        end

        subgraph EngineSpider[" "]
            SubmitS["Engine._submit_spider()"]
            RouteS["Engine._route_spider_value()"]
        end

        subgraph EnginePipeline[" "]
            SubmitP["Engine._submit_pipeline()"]
            RouteP["Engine._route_pipeline_value()"]
        end

        HandleEvent["Engine._handle_engine_event()"]
        Check["Engine._check_finished()"]
    end

    subgraph Components[" "]
        subgraph DownloaderModule["Downloader"]
            DIn["DownloaderMiddleware.process_input()"]
            Downloader["Downloader.fetch()"]
            DOut["DownloaderMiddleware.process_output()"]
        end

        subgraph SpiderModule["Spider"]
            SIn["SpiderMiddleware.process_input()"]
            Spider["Spider.parse()"]
            SOut["SpiderMiddleware.process_output()"]
        end

        subgraph PipelineModule["Pipeline"]
            PIn["PipelineMiddleware.process_input()"]
            Pipeline["Pipeline.process_item()"]
            POut["PipelineMiddleware.process_output()"]
        end
    end


    Start["Spider.start_requests()"] -->|Request| Seed
    Seed -->|Request| Enqueue
    Enqueue --> Dequeue
    Dequeue -->|Request| SubmitD

    SubmitD --> DIn
    DIn --> Downloader
    Downloader --> DOut
    DOut --> RouteD

    RouteD -->|Request| SubmitD
    RouteD -->|Response| SubmitS
    RouteD -->|Failure| Check
    RouteD -->|EngineEvent| HandleEvent

    SubmitS --> SIn
    SIn --> Spider
    Spider --> SOut
    SOut --> RouteS

    RouteS -->|Request| Enqueue
    RouteS -->|Item| SubmitP
    RouteS -->|Failure| Check
    RouteS -->|EngineEvent| HandleEvent

    SubmitP --> PIn
    PIn --> Pipeline
    Pipeline --> POut
    POut --> RouteP

    RouteP -->|Failure| Check
    RouteP -->|EngineEvent| HandleEvent

    Check --> Done
    Check --> Failed
Loading

Model Types

A WorkOrder is the engine's internal unit of work for one leased root Request. It:

  • Tracks all downloader, spider, and pipeline tasks spawned from that request.
  • Collects per-order payload and counters.
  • Is marked done or failed only after all related tasks are idle.
Type Fields Purpose
Request url, method, headers, body, meta, priority A crawl request. Spiders create it, schedulers queue it, downloaders fetch it. Duplicate filtering uses fingerprint(), which is based on method, URL, and body.
Response url, status, body, headers, request, cookies Downloader output and spider input. Use text for charset-aware body decoding.
Item data A mapping-like object emitted by spiders and consumed by pipelines.
Failure value, exception, stage A component-level failure. Returning it marks the current work order as failed.
EngineEvent action, payload, log A control message from a component to the engine. Details are below.

EngineEvent Actions:

Action Payload Effect
EngineAction.NONE Optional Does not change runtime state. If log is set, the engine emits that log record.
EngineAction.COLLECT Any mapping Merges values into the current WorkOrder.payload. Integer values are accumulated as counters; other values replace the previous value for that key.
EngineAction.STOP_GRACEFULLY Optional reason Stops taking new scheduler requests, lets running work finish, then emits a stopped log record.
EngineAction.STOP_NOW Optional reason Fails the current work order, cancels running futures where possible, stops scheduling immediately, then emits a stopped log record.

Example:

from rubberneck import EngineAction, EngineEvent

yield EngineEvent(EngineAction.COLLECT, {"pages": 1, "last_url": response.url})

After several responses, pages is accumulated while last_url keeps the most recent URL.

Minimal Crawler

from rubberneck import Engine, Item, Request, Response, Spider


class ExampleSpider(Spider):
    name = "example"

    def start_requests(self):
        yield Request("https://example.org/")

    def parse(self, response: Response):
        yield Item({"url": response.url, "status": response.status})


Engine(ExampleSpider()).run()

Components

Runtime components can be passed as an instance, a registry name, or a ComponentSpec with constructor options.

from rubberneck import ComponentSpec, Engine

engine = Engine(
    spider,
    scheduler='sqlite',
    downloader='session_pool',
    pipeline=ComponentSpec('sqlite', {'table': 'items'}),
    logger='standard',
)
Component Key method Input and output Provided implementations
Engine run() Input: a Spider plus optional scheduler, downloader, pipeline, logger, middleware lists, and worker counts. Output: EngineStats. Engine coordinates lifecycle, worker pools, routing, stop handling, and final done/failed acknowledgement.
Scheduler enqueue(request) Input: a Request. Output: True if accepted, False if filtered as duplicate. memory is an in-memory queue for tests and short runs. sqlite is durable and recovers leased/failed requests as pending on restart.
dequeue() Input: none. Output: the next leased Request, or None when no request is pending.
mark_done(request) Input: a leased Request. Output: none; records successful completion.
mark_failed(request, error) Input: a leased Request and an exception. Output: none; records failed completion.
Downloader fetch(request) Input: one Request. Output: Response, downloader-local Request, Failure, or EngineEvent. session_pool uses reusable requests.Session objects. urllib uses the standard library.
DownloaderMiddleware process_input(request) Input: a Request before fetch. Output: the request to pass to the downloader. cookies manages cookie jars. referer sets and propagates Referer. RetryDownloaderMiddleware retries failures. ChallengeDownloaderMiddleware is a base class for challenge flows.
process_output(request, output) Input: the request and downloader output stream. Output: modified output stream.
Spider start_requests() Input: none. Output: seed Request values for the scheduler. No built-in implementation is provided yet; implement it for each site.
parse(response) Input: one Response. Output: Request, Item, Failure, or EngineEvent.
SpiderMiddleware process_input(response) Input: a Response before parsing. Output: the response to pass to the spider. No built-in implementation is provided yet.
process_output(response, output) Input: the response and spider output stream. Output: modified output stream.
Pipeline process_item(item) Input: one Item. Output: Failure or EngineEvent. A successful call increments the per-order processed counter. sqlite writes item data to SQLite, creating the table and adding columns as item keys appear.
PipelineMiddleware process_input(item) Input: an Item before pipeline processing. Output: the item to pass to the pipeline. No built-in implementation is provided yet.
process_output(item, output) Input: the item and pipeline output stream. Output: modified output stream.
Logger emit(record) Input: one LogRecord. Output: none. standard emits records through Python logging, with action filtering and periodic summaries.

Installation

From the project root directory, install the package with:

python -m pip install .

For development, install it in editable mode:

python -m pip install -e .

About

A web crawling framework.

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages