Pipelines

Pipelines process scraped items and write them to files, databases, or external services. They are executed in order and each pipeline receives the output of the previous one. See src/silkworm/pipelines.py.

Pipeline Interface

Every pipeline implements three core async methods:

class ItemPipeline:
    async def open(self, spider) -> None: ...
    async def close(self, spider) -> None: ...
    async def process_item(self, item, spider): ...

class BatchItemPipeline(ItemPipeline):
    async def process_items(self, items, spider): ...

All built-in pipelines implement BatchItemPipeline. By default the engine calls process_item() immediately for every emitted item. Set item_batch_size above one to let the engine combine emitted items and use native bulk operations; partial batches flush after item_batch_wait (0.05 seconds by default), when a callback drains, and at shutdown.

run_spider(
    MySpider,
    item_pipelines=[SQLitePipeline("data/items.db")],
    item_batch_size=100,
    item_batch_wait=0.05,
)

Database, service, streaming-file, and buffered built-ins opt into native batching. Callback and validation stages remain per-item unless CallbackPipeline receives a batch_callback, preserving transformations and DropItem statistics. Custom pipelines may set native_batch = True when process_items() returns exactly one ordered result for every input and never uses DropItem for individual rows.

Explicit process_items() calls remain available. A client bulk request can partially persist data before reporting a failure; any reported row failure raises for the batch.

Pipeline Usage

from silkworm.pipelines import JsonLinesPipeline, SQLitePipeline

run_spider(
    MySpider,
    item_pipelines=[
        JsonLinesPipeline("data/items.jl"),
        SQLitePipeline("data/items.db", table="items"),
    ],
)

Streaming vs Buffered Pipelines

  • Streaming (per item): JsonLinesPipeline, CSVPipeline, XMLPipeline, SQLitePipeline, WebhookPipeline (batch_size=1), ZenohPipeline.

  • Buffered (write on close): PolarsPipeline, ExcelPipeline, YAMLPipeline, AvroPipeline, VortexPipeline, S3JsonLinesPipeline, FTPPipeline, SFTPPipeline, RssPipeline.

  • Native bulk ingestion: SQL/database pipelines, Iggy, Elasticsearch, Google Sheets, Webhook, Taskiq, Zenoh, and coalesced file writers.

Note: Buffered pipelines keep items in memory. Prefer streaming ones for large crawls.

Optional Dependencies

Pipelines backed by optional extras are always importable; constructing one without its extra installed raises ImportError with the install command. Each module also exports an availability flag you can check first:

Flag

Pipelines

FASTAVRO_AVAILABLE

AvroPipeline

CASSANDRA_AVAILABLE

CassandraPipeline

AIOCOUCH_AVAILABLE

CouchDBPipeline

DUCKDB_AVAILABLE

DuckDBPipeline

AIOBOTO3_AVAILABLE

DynamoDBPipeline

ELASTICSEARCH_AVAILABLE

ElasticsearchPipeline

OPENPYXL_AVAILABLE

ExcelPipeline

AIOFTP_AVAILABLE

FTPPipeline

GOOGLE_SHEETS_AVAILABLE

GoogleSheetsPipeline

IGGY_AVAILABLE

IggyPipeline

MOTOR_AVAILABLE

MongoDBPipeline

ORMSGPACK_AVAILABLE

MsgPackPipeline

AIOMYSQL_AVAILABLE

MySQLPipeline

POLARS_AVAILABLE

PolarsPipeline

ASYNCPG_AVAILABLE

PostgreSQLPipeline

OPENDAL_AVAILABLE

S3JsonLinesPipeline, JsonLinesPipeline(use_opendal=True)

ASYNCSSH_AVAILABLE

SFTPPipeline

SNOWFLAKE_AVAILABLE

SnowflakePipeline

TASKIQ_AVAILABLE

TaskiqPipeline

VORTEX_AVAILABLE

VortexPipeline

WREQ_AVAILABLE

WebhookPipeline

YAML_AVAILABLE

YAMLPipeline

ZENOH_AVAILABLE

ZenohPipeline

from silkworm.pipelines import POLARS_AVAILABLE, JsonLinesPipeline, PolarsPipeline

pipeline = PolarsPipeline("data/items.parquet") if POLARS_AVAILABLE else JsonLinesPipeline("data/items.jl")

Per-item Logging

Wrap any pipeline in LoggedPipeline(pipeline, log_level=...) to change or silence (log_level=None) its per-item log messages. See Logging and Stats.

Dropping Items

Raise DropItem(message, reason=...) from process_item to discard an item: later pipelines are skipped, the item is counted under items_dropped (labeled by reason), and emit() returns normally. Only items that pass every pipeline count as items_scraped.

Built-in Pipelines

ValidationPipeline

Validates items against a Pydantic model (anything with model_validate) or a validator function, forwarding normalized items and dropping invalid ones with reason invalid (or raising with on_invalid="raise"). The first log_limit failures are logged with their validation errors. Combine it with the engine’s max_item_drop_rate or min_items to fail crawls when extraction breaks. See Validating items.

from pydantic import BaseModel
from silkworm.pipelines import JsonLinesPipeline, ValidationPipeline


class Quote(BaseModel):
    text: str
    author: str


run_spider(
    QuotesSpider,
    item_pipelines=[ValidationPipeline(Quote), JsonLinesPipeline("quotes.jl")],
    max_item_drop_rate=0.05,
)

CallbackPipeline

  • Purpose: Run a custom callback for each item (sync or async).

  • Behavior: If the callback returns None, the original item passes through unchanged.

  • Options: callback (an ItemCallback), log_level for per-item log messages.

  • Extras: none.

  • Code: src/silkworm/_pipelines/callback_pipeline.py

from silkworm.pipelines import CallbackPipeline

async def validate_item(item, spider):
    return item

CallbackPipeline(callback=validate_item)

JsonLinesPipeline

  • Purpose: Write items as JSON Lines to a local file.

  • Options: path, use_opendal (async writes with OpenDAL when available), log_level for per-item log messages.

  • Extras: s3 (OpenDAL).

  • Code: src/silkworm/_pipelines/jsonlines_pipeline.py

JsonLinesPipeline("data/items.jl", use_opendal=False)

MsgPackPipeline

MsgPackPipeline("data/items.msgpack", mode="append")

SQLitePipeline

SQLitePipeline("data/items.db", table="quotes")

XMLPipeline

XMLPipeline("data/items.xml", root_element="items", item_element="item")

RssPipeline

  • Purpose: Write items to an RSS 2.0 feed (buffered).

  • Options: path, channel_title, channel_link, channel_description, max_items (most recent items kept, default 50; None for no limit).

  • Field mappings: item_title_field (default "title"), item_link_field ("link"), item_description_field ("description"), and optional item_pub_date_field, item_guid_field, item_author_field.

  • Extras: none.

  • Code: src/silkworm/_pipelines/rss_pipeline.py

RssPipeline(
    "data/feed.xml",
    channel_title="My Feed",
    channel_link="https://example.com",
    channel_description="Latest items",
    max_items=50,
)

CSVPipeline

CSVPipeline("data/items.csv", fieldnames=["author", "text", "tags"])

TaskiqPipeline

TaskiqPipeline(broker, task_name=".:process_item")

ZenohPipeline

  • Purpose: Publish JSON-serialized items to Zenoh immediately.

  • Options: key_expr (static string or a sync/async ZenohKeyResolver taking (item, spider)), config or session (mutually exclusive), encoding (default "application/json"), plus Zenoh publisher QoS options congestion_control, priority, express, reliability, and allowed_destination.

  • Lifecycle: Opens and closes its own session by default. An injected session remains caller-owned. All publishers created by the pipeline are undeclared on close.

  • Routing: Dynamic publishers are declared lazily and cached by key expression until close.

  • Extras: zenoh.

  • Code: src/silkworm/_pipelines/zenoh_pipeline.py

from silkworm.pipelines import ZenohPipeline

ZenohPipeline("scraping/items")

IggyPipeline

  • Purpose: Publish each item as a compact JSON message to Apache Iggy.

  • Options: stream, topic, optional connection_string or connected client, partitioning, producer mode, stream/topic creation settings, and retry settings.

  • Behavior: Uses Iggy’s high-level producer. process_item(...) publishes one message, while process_items(...) publishes a list as one native Iggy batch. Background producer modes are flushed during pipeline shutdown.

  • Extras: iggy.

  • Code: src/silkworm/_pipelines/iggy_pipeline.py

from apache_iggy import BackgroundProducerConfig
from silkworm.pipelines import IggyPipeline

pipeline = IggyPipeline(
    "scraping",
    "items",
    connection_string="iggy+tcp://iggy:iggy@127.0.0.1:8090",
    mode=BackgroundProducerConfig(),
)

await pipeline.open(spider)
try:
    await pipeline.process_items([{"id": 1}, {"id": 2}], spider)
finally:
    await pipeline.close(spider)

PolarsPipeline

PolarsPipeline("data/items.parquet", mode="append")

ExcelPipeline

ExcelPipeline("data/items.xlsx", sheet_name="quotes")

YAMLPipeline

YAMLPipeline("data/items.yaml")

AvroPipeline

AvroPipeline("data/items.avro", schema=my_schema)

ElasticsearchPipeline

ElasticsearchPipeline(hosts=["http://localhost:9200"], index="quotes")

MongoDBPipeline

MongoDBPipeline(database="scraping", collection="items")

S3JsonLinesPipeline

  • Purpose: Write JSON Lines to S3 via OpenDAL (buffered).

  • Options: bucket, key, region, optional endpoint, access_key_id, secret_access_key.

  • Extras: s3.

  • Code: src/silkworm/_pipelines/s3_pipeline.py

S3JsonLinesPipeline(bucket="my-bucket", key="data/items.jl")

VortexPipeline

VortexPipeline("data/items.vortex")

MySQLPipeline

MySQLPipeline(database="scraping", table="items")

PostgreSQLPipeline

PostgreSQLPipeline(database="scraping", table="items")

WebhookPipeline

  • Purpose: Send items to a webhook using wreq.

  • Options: url, method, headers, timeout, batch_size.

  • Behavior: If batch_size > 1, the payload is a list of items.

  • Extras: none (wreq is core).

  • Code: src/silkworm/_pipelines/webhook_pipeline.py

WebhookPipeline("https://example.com/webhook", batch_size=10)

GoogleSheetsPipeline

GoogleSheetsPipeline(
    spreadsheet_id="...",
    credentials_file="creds.json",
    sheet_name="items",
    batch_size=100,
)

SnowflakePipeline

SnowflakePipeline(
    account="acct",
    user="user",
    password="pass",
    database="db",
    schema="PUBLIC",
    warehouse="WH",
    table="items",
)

FTPPipeline

FTPPipeline(host="ftp.example.com", user="user", password="pass")

SFTPPipeline

  • Purpose: Upload JSON Lines to SFTP (buffered).

  • Options: host, user, password or private_key, remote_path, port, known_hosts, verify_host_key.

  • Host key verification: on by default. The server’s host key is checked against ~/.ssh/known_hosts; pass known_hosts="path/to/known_hosts" to use another file. verify_host_key=False skips the check. Use it only on trusted networks, since it allows man-in-the-middle attacks. It can’t be combined with known_hosts.

  • Extras: sftp.

  • Code: src/silkworm/_pipelines/sftp_pipeline.py

SFTPPipeline(host="sftp.example.com", user="user", password="pass")

# Verify against a specific known_hosts file
SFTPPipeline(
    host="sftp.example.com",
    user="user",
    private_key="~/.ssh/id_ed25519",
    known_hosts="deploy/known_hosts",
)

Behaviour change: earlier versions never verified the SFTP server’s host key. Uploads to a server that isn’t in your known_hosts file now fail until you add its key (for example with ssh-keyscan sftp.example.com >> ~/.ssh/known_hosts), point known_hosts at a file that has it, or opt out with verify_host_key=False.

CassandraPipeline

CassandraPipeline(hosts=["127.0.0.1"], keyspace="scraping", table="items")

CouchDBPipeline

CouchDBPipeline(url="http://localhost:5984", database="scraping")

DynamoDBPipeline

  • Purpose: Insert items into DynamoDB (auto-creates table if missing).

  • Options: table_name, region_name, aws_access_key_id, aws_secret_access_key, endpoint_url.

  • Extras: dynamodb.

  • Code: src/silkworm/_pipelines/dynamodb_pipeline.py

DynamoDBPipeline(table_name="items", region_name="us-east-1")

DuckDBPipeline

DuckDBPipeline(database="data/items.db", table="items")