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 |
|---|---|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
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(anItemCallback),log_levelfor per-item log messages.Extras: none.
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_levelfor per-item log messages.Extras:
s3(OpenDAL).
JsonLinesPipeline("data/items.jl", use_opendal=False)
MsgPackPipeline¶
Purpose: Binary MessagePack file using
ormsgpack.Options:
path,mode(writeorappend).Extras:
msgpack.
MsgPackPipeline("data/items.msgpack", mode="append")
SQLitePipeline¶
Purpose: Store items as JSON text in SQLite.
Options:
path,table.Extras: none.
SQLitePipeline("data/items.db", table="quotes")
XMLPipeline¶
Purpose: Write items as XML with nested data preserved.
Options:
path,root_element,item_element.Extras: none.
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;Nonefor no limit).Field mappings:
item_title_field(default"title"),item_link_field("link"),item_description_field("description"), and optionalitem_pub_date_field,item_guid_field,item_author_field.Extras: none.
RssPipeline(
"data/feed.xml",
channel_title="My Feed",
channel_link="https://example.com",
channel_description="Latest items",
max_items=50,
)
CSVPipeline¶
Purpose: CSV export (nested dicts flattened, lists joined by commas).
Options:
path,fieldnames(optional).Extras: none.
CSVPipeline("data/items.csv", fieldnames=["author", "text", "tags"])
TaskiqPipeline¶
Purpose: Send items to a Taskiq broker/queue.
Options:
broker,taskortask_name.Extras:
taskiq.
TaskiqPipeline(broker, task_name=".:process_item")
ZenohPipeline¶
Purpose: Publish JSON-serialized items to Zenoh immediately.
Options:
key_expr(static string or a sync/asyncZenohKeyResolvertaking(item, spider)),configorsession(mutually exclusive),encoding(default"application/json"), plus Zenoh publisher QoS optionscongestion_control,priority,express,reliability, andallowed_destination.Lifecycle: Opens and closes its own session by default. An injected
sessionremains 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.
from silkworm.pipelines import ZenohPipeline
ZenohPipeline("scraping/items")
IggyPipeline¶
Purpose: Publish each item as a compact JSON message to Apache Iggy.
Options:
stream,topic, optionalconnection_stringor connectedclient,partitioning, producermode, stream/topic creation settings, and retry settings.Behavior: Uses Iggy’s high-level producer.
process_item(...)publishes one message, whileprocess_items(...)publishes a list as one native Iggy batch. Background producer modes are flushed during pipeline shutdown.Extras:
iggy.
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¶
Purpose: Write Parquet via Polars (buffered).
Options:
path,mode(writeorappend).Extras:
polars.
PolarsPipeline("data/items.parquet", mode="append")
ExcelPipeline¶
Purpose: Write XLSX via openpyxl (buffered, flattening like CSV).
Options:
path,sheet_name.Extras:
excel.
ExcelPipeline("data/items.xlsx", sheet_name="quotes")
YAMLPipeline¶
Purpose: Write YAML (buffered).
Options:
path.Extras:
yaml.
YAMLPipeline("data/items.yaml")
AvroPipeline¶
Purpose: Write Avro (buffered). Schema can be inferred.
Options:
path,schema(optional).Extras:
avro.
AvroPipeline("data/items.avro", schema=my_schema)
ElasticsearchPipeline¶
Purpose: Index items in Elasticsearch.
Options:
hosts,index,**es_kwargs.Extras:
elasticsearch.
ElasticsearchPipeline(hosts=["http://localhost:9200"], index="quotes")
MongoDBPipeline¶
Purpose: Insert items into MongoDB.
Options:
connection_string,database,collection.Extras:
mongodb.
MongoDBPipeline(database="scraping", collection="items")
S3JsonLinesPipeline¶
Purpose: Write JSON Lines to S3 via OpenDAL (buffered).
Options:
bucket,key,region, optionalendpoint,access_key_id,secret_access_key.Extras:
s3.
S3JsonLinesPipeline(bucket="my-bucket", key="data/items.jl")
VortexPipeline¶
Purpose: Write Vortex columnar format (buffered).
Options:
path.Extras:
vortex.
VortexPipeline("data/items.vortex")
MySQLPipeline¶
Purpose: Insert items into MySQL as JSON.
Options:
host,port,user,password,database,table.Extras:
mysql.
MySQLPipeline(database="scraping", table="items")
PostgreSQLPipeline¶
Purpose: Insert items into PostgreSQL as JSONB.
Options:
host,port,user,password,database,table.Extras:
postgresql.
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).
WebhookPipeline("https://example.com/webhook", batch_size=10)
GoogleSheetsPipeline¶
Purpose: Append rows to Google Sheets (batching, flattening like CSV).
Options:
spreadsheet_id,credentials_file,sheet_name,batch_size.Extras:
gsheets.
GoogleSheetsPipeline(
spreadsheet_id="...",
credentials_file="creds.json",
sheet_name="items",
batch_size=100,
)
SnowflakePipeline¶
Purpose: Insert items into Snowflake as JSON.
Options:
account,user,password,database,schema,warehouse,table,role.Extras:
snowflake.
SnowflakePipeline(
account="acct",
user="user",
password="pass",
database="db",
schema="PUBLIC",
warehouse="WH",
table="items",
)
FTPPipeline¶
Purpose: Upload JSON Lines to FTP (buffered).
Options:
host,user,password,remote_path,port.Extras:
ftp.
FTPPipeline(host="ftp.example.com", user="user", password="pass")
SFTPPipeline¶
Purpose: Upload JSON Lines to SFTP (buffered).
Options:
host,user,passwordorprivate_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; passknown_hosts="path/to/known_hosts"to use another file.verify_host_key=Falseskips the check. Use it only on trusted networks, since it allows man-in-the-middle attacks. It can’t be combined withknown_hosts.Extras:
sftp.
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_hostsfile now fail until you add its key (for example withssh-keyscan sftp.example.com >> ~/.ssh/known_hosts), pointknown_hostsat a file that has it, or opt out withverify_host_key=False.
CassandraPipeline¶
Purpose: Insert items into Cassandra.
Options:
hosts,keyspace,table,username,password,port.Extras:
cassandra(not available on Windows).
CassandraPipeline(hosts=["127.0.0.1"], keyspace="scraping", table="items")
CouchDBPipeline¶
Purpose: Insert items into CouchDB.
Options:
url,database,username,password.Extras:
couchdb.
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.
DynamoDBPipeline(table_name="items", region_name="us-east-1")
DuckDBPipeline¶
Purpose: Insert items into DuckDB as JSON.
Options:
database,table.Extras:
duckdb.
DuckDBPipeline(database="data/items.db", table="items")