silkworm.pipelines

Item pipelines.

Implementations live in silkworm._pipelines; this module re-exports the public API.

class silkworm.pipelines.AvroPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to an Avro file.

Avro is a row-oriented data serialization system with compact binary format. Items are buffered until the pipeline closes.

Parameters:
  • path – Output Avro file.

  • schema – Explicit Avro schema. When omitted, a simple record schema is inferred from the first item.

Example:

from silkworm.pipelines import AvroPipeline

schema = {
    "type": "record",
    "name": "Quote",
    "fields": [
        {"name": "text", "type": "string"},
        {"name": "author", "type": "string"},
        {"name": "tags", "type": {"type": "array", "items": "string"}},
    ],
}
pipeline = AvroPipeline("data/items.avro", schema=schema)
__init__(path='items.avro', *, schema=None)[source]

Initialize AvroPipeline.

Parameters:
  • path (str | Path) – Path to the output file (default: “items.avro”)

  • schema (Mapping[str, JSONLike] | None) – Avro schema dict. If None, will infer from first item.

Return type:

None

async open(spider)[source]

Create parent directories and reset the in-memory item buffer.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Infer or apply the schema and write all buffered items.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Buffer one item for serialization during close().

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

type silkworm.pipelines.BatchItemCallback = Callable[[list[JSONValue], Spider], list[JSONValue] | Awaitable[list[JSONValue] | None] | None]
class silkworm.pipelines.BatchItemPipeline[source]

Bases: ItemPipeline, Protocol

Item pipeline that can process an explicit batch in one call.

async process_items(items, spider)[source]

Process a batch in order and return items for the next pipeline.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.CSVPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Stream flattened items to a UTF-8 CSV file.

Parameters:
  • path – Destination file path.

  • fieldnames – Fixed column order. When omitted, columns are inferred from the first item.

Nested mappings use underscore-separated keys and list values are joined with commas.

__init__(path='items.csv', *, fieldnames=None)[source]
Parameters:
Return type:

None

async open(spider)[source]

Create parent directories and open a fresh CSV destination.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Flush and close the CSV file if it is open.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Flatten and append one mapping item to the CSV file.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.CallbackPipeline[source]

Bases: _BatchPipelineMixin

Pipeline that invokes a callback function to process each item.

This pipeline allows you to define custom item processing logic using a simple callback function, making it easy to handle items without creating a full pipeline class.

The callback function can be either synchronous or asynchronous and receives the item and spider as arguments.

Parameters:
  • callback – Synchronous or asynchronous callable receiving (item, spider). Returning None preserves the original item; any other result is forwarded to the next pipeline.

  • log_level – Severity used for per-item processing logs.

Example:

from silkworm.pipelines import CallbackPipeline

def process_item(item, spider):
    # Your custom processing logic
    print(f"Processing item from {spider.name}: {item}")
    return item

pipeline = CallbackPipeline(callback=process_item)

# Or with an async callback:
async def async_process_item(item, spider):
    # Your async processing logic
    await some_async_operation(item)
    return item

pipeline = CallbackPipeline(callback=async_process_item)
__init__(callback, *, batch_callback=None, log_level='DEBUG')[source]

Initialize CallbackPipeline.

Parameters:
  • callback (ItemCallback) – A callable that takes (item, spider) and returns the processed item. Can be either synchronous or asynchronous.

  • batch_callback (BatchItemCallback | None)

  • log_level (LogLevel)

Return type:

None

async open(spider)[source]

Open the pipeline.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Close the pipeline.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Process an item using the callback function.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.CassandraPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to an Apache Cassandra database.

Parameters:
  • hosts – Cassandra cluster hosts.

  • keyspace – Keyspace created with a single-node replication strategy when absent.

  • table – Table storing UUID, spider, JSON data, and creation timestamp.

  • username – Optional authentication username.

  • password – Optional authentication password.

  • port – Cassandra native protocol port.

Example:

from silkworm.pipelines import CassandraPipeline

pipeline = CassandraPipeline(
    hosts=["127.0.0.1"],
    keyspace="scraping",
    table="items",
    username="cassandra",
    password="cassandra",
)
__init__(hosts=None, keyspace='scraping', *, table='items', username=None, password=None, port=9042)[source]

Initialize CassandraPipeline.

Parameters:
  • hosts (list[str] | None) – List of Cassandra cluster hosts (default: [“127.0.0.1”])

  • keyspace (str) – Keyspace name

  • table (str) – Table name (default: “items”)

  • username (str | None) – Optional username for authentication

  • password (str | None) – Optional password for authentication

  • port (int) – Cassandra port (default: 9042)

Return type:

None

async open(spider)[source]

Connect and create the keyspace and JSON-document table if absent.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Shut down the Cassandra cluster connection.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Insert one item with a UUID, spider name, and timestamp.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.CouchDBPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to a CouchDB database.

Parameters:
  • url – CouchDB server URL.

  • database – Database opened or created during startup.

  • username – Optional authentication username.

  • password – Optional authentication password.

Example:

from silkworm.pipelines import CouchDBPipeline

pipeline = CouchDBPipeline(
    url="http://localhost:5984",
    database="scraping",
    username="admin",
    password="password",
)
__init__(url='http://localhost:5984', database='scraping', *, username=None, password=None)[source]

Initialize CouchDBPipeline.

Parameters:
  • url (str) – CouchDB server URL (default: “http://localhost:5984”)

  • database (str) – Database name (default: “scraping”)

  • username (str | None) – Optional username for authentication

  • password (str | None) – Optional password for authentication

Return type:

None

async open(spider)[source]

Connect and open or create the configured CouchDB database.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Exit the CouchDB client context and release database references.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Create a document containing the spider name and original item.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.DuckDBPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to a DuckDB database.

DuckDB is an embedded analytical database with excellent performance for OLAP queries. This pipeline stores items in a DuckDB table as JSON.

Parameters:
  • database – Embedded DuckDB database file.

  • table – Valid unquoted destination table name.

Example:

from silkworm.pipelines import DuckDBPipeline

pipeline = DuckDBPipeline(
    database="data/scraping.db",
    table="items",
)
__init__(database='items.db', *, table='items')[source]

Initialize DuckDBPipeline.

Parameters:
  • database (str | Path) – Path to the DuckDB database file (default: “items.db”)

  • table (str) – Table name (default: “items”)

Return type:

None

async open(spider)[source]

Open DuckDB and create the sequence and JSON table if absent.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Close the embedded DuckDB connection.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Insert one JSON item with its spider name.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.DynamoDBPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to AWS DynamoDB.

Parameters:
  • table_name – Table opened or created with a string id hash key.

  • region_name – AWS region.

  • aws_access_key_id – Explicit access key, or None for normal provider discovery.

  • aws_secret_access_key – Explicit secret paired with the access key.

  • endpoint_url – Custom endpoint for DynamoDB Local or compatible services.

Example:

from silkworm.pipelines import DynamoDBPipeline

pipeline = DynamoDBPipeline(
    table_name="items",
    region_name="us-east-1",
    aws_access_key_id="YOUR_KEY",
    aws_secret_access_key="YOUR_SECRET",
)
__init__(table_name='items', *, region_name='us-east-1', aws_access_key_id=None, aws_secret_access_key=None, endpoint_url=None)[source]

Initialize DynamoDBPipeline.

Parameters:
  • table_name (str) – DynamoDB table name (default: “items”)

  • region_name (str) – AWS region (default: “us-east-1”)

  • aws_access_key_id (str | None) – AWS access key ID (uses env vars/IAM if not provided)

  • aws_secret_access_key (str | None) – AWS secret access key (uses env vars/IAM if not provided)

  • endpoint_url (str | None) – Custom endpoint URL for DynamoDB Local or other services

Return type:

None

async open(spider)[source]

Open DynamoDB resources and create the keyed table if absent.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Exit the DynamoDB client and resource contexts.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Store one item with a generated string ID and spider metadata.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.ElasticsearchPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to an Elasticsearch index.

Parameters:
  • hosts – One endpoint or a list of Elasticsearch endpoints.

  • index – Destination index name.

  • **es_kwargs – Additional AsyncElasticsearch client options.

Example:

from silkworm.pipelines import ElasticsearchPipeline

pipeline = ElasticsearchPipeline(
    hosts=["http://localhost:9200"],
    index="quotes",
)
__init__(hosts='http://localhost:9200', *, index='items', **es_kwargs)[source]

Initialize ElasticsearchPipeline.

Parameters:
  • hosts (list[str] | str) – Elasticsearch host(s)

  • index (str) – Index name

  • **es_kwargs (Any) – Additional kwargs for AsyncElasticsearch client

Return type:

None

async open(spider)[source]

Create the asynchronous Elasticsearch client.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Close the Elasticsearch transport and release the client.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Index one item as a document in the configured index.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.ExcelPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to an Excel file (.xlsx).

Items are buffered, flattened, and written when the pipeline closes.

Parameters:
  • path – Output workbook path.

  • sheet_name – Worksheet title.

Example:

from silkworm.pipelines import ExcelPipeline

pipeline = ExcelPipeline("data/items.xlsx", sheet_name="quotes")
__init__(path='items.xlsx', *, sheet_name='Sheet1')[source]

Initialize ExcelPipeline.

Parameters:
  • path (str | Path) – Path to the output file (default: “items.xlsx”)

  • sheet_name (str) – Name of the Excel sheet (default: “Sheet1”)

Return type:

None

async open(spider)[source]

Create parent directories and reset the workbook item buffer.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Flatten buffered items and write the XLSX workbook.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Buffer one item for workbook creation during close().

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.FTPPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to an FTP server in JSON Lines format.

Parameters:
  • host – FTP server hostname.

  • user – Login username.

  • password – Login password.

  • remote_path – Destination object path.

  • port – FTP control port.

Items are buffered locally and replace the remote file during close.

Example:

from silkworm.pipelines import FTPPipeline

pipeline = FTPPipeline(
    host="ftp.example.com",
    user="username",
    password="password",
    remote_path="data/items.jl",
)
__init__(host, user, password, remote_path='items.jl', *, port=21)[source]

Initialize FTPPipeline.

Parameters:
  • host (str) – FTP server hostname

  • user (str) – FTP username

  • password (str) – FTP password

  • remote_path (str) – Remote file path (default: “items.jl”)

  • port (int) – FTP port (default: 21)

Return type:

None

async open(spider)[source]

Reset the in-memory JSON Lines buffer without connecting yet.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Upload all buffered lines and close the FTP connection.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Serialize and buffer one JSON line for upload during close().

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.GoogleSheetsPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that appends items to a Google Sheet.

Requires Google Sheets API credentials (service account JSON file).

Parameters:
  • spreadsheet_id – Spreadsheet identifier from its URL.

  • credentials_file – Service-account credentials JSON file.

  • sheet_name – Worksheet receiving rows.

  • batch_size – Items flattened and appended per API batch.

Example:

from silkworm.pipelines import GoogleSheetsPipeline

pipeline = GoogleSheetsPipeline(
    spreadsheet_id="1BxiMVs0XRA5nFMdKvBdBZjgmUUqptlbs74OgvE2upms",
    credentials_file="path/to/credentials.json",
    sheet_name="Sheet1",
)
__init__(spreadsheet_id, credentials_file, *, sheet_name='Sheet1', batch_size=100)[source]

Initialize GoogleSheetsPipeline.

Parameters:
  • spreadsheet_id (str) – Google Sheets spreadsheet ID (from the URL)

  • credentials_file (str) – Path to service account credentials JSON file

  • sheet_name (str) – Name of the sheet to append to (default: “Sheet1”)

  • batch_size (int) – Number of items to batch before writing (default: 100)

Return type:

None

async open(spider)[source]

Authenticate the Sheets service and reset batching state.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Write any partial batch and release the Sheets service reference.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Buffer one item and append a batch when batch_size is reached.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.IggyPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that publishes JSON-serialized items to Apache Iggy.

Parameters:
  • stream – Name of the destination Iggy stream.

  • topic – Name of the destination topic.

  • connection_string – Iggy connection string used to create a pipeline-owned client. Defaults to a local TCP server with Iggy’s development credentials. Mutually exclusive with client.

  • client – Existing connected Iggy client to reuse. The pipeline does not connect or otherwise manage an injected client.

  • partitioning – Iggy partitioning strategy. Balanced partitioning is used when omitted.

  • mode – Direct or background producer configuration. Iggy’s direct mode is used when omitted.

  • create_stream_if_not_exists – Create stream during startup when needed.

  • create_topic_if_not_exists – Create topic during startup when needed.

  • topic_partitions_count – Number of partitions used when creating topic.

  • send_retries – Number of producer send retries, or None for unlimited retries.

  • send_retry_interval – Delay between producer send retries.

Each item is encoded as one compact UTF-8 JSON message. Closing the pipeline shuts down its producer and flushes messages accepted by a background producer. Use process_items() to publish several items in one native Iggy batch.

Example:

from silkworm.pipelines import IggyPipeline

pipeline = IggyPipeline("scraping", "items")
__init__(stream, topic, *, connection_string=None, client=None, partitioning=None, mode=None, create_stream_if_not_exists=True, create_topic_if_not_exists=True, topic_partitions_count=1, send_retries=3, send_retry_interval=timedelta(seconds=1))[source]

Initialize an Apache Iggy producer pipeline.

Parameters:
  • stream (str)

  • topic (str)

  • connection_string (str | None)

  • client (_IggyClient | None)

  • partitioning (_Partitioning | None)

  • mode (DirectProducerConfig | BackgroundProducerConfig | None)

  • create_stream_if_not_exists (bool)

  • create_topic_if_not_exists (bool)

  • topic_partitions_count (int)

  • send_retries (int | None)

  • send_retry_interval (timedelta)

Return type:

None

async open(spider)[source]

Connect a pipeline-owned client and initialize the Iggy producer.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Flush and shut down the Iggy producer.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Publish one JSON item to Iggy and pass it to the next pipeline.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Publish several JSON items in one Iggy batch and return them unchanged.

Parameters:
Return type:

list[JSONValue]

type silkworm.pipelines.ItemCallback = Callable[[JSONValue, Spider], JSONValue | Awaitable[JSONValue | None] | None]
class silkworm.pipelines.ItemPipeline[source]

Bases: Protocol

Protocol implemented by ordered item-processing stages.

async open(spider)[source]

Allocate resources once before the first item is processed.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Flush buffered data and release resources after the crawl.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Process and return the item passed to the next pipeline.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

__init__(*args, **kwargs)
type silkworm.pipelines.ItemSchema = ModelSchema | type[object] | ItemValidator
type silkworm.pipelines.ItemValidator = Callable[[JSONValue], JSONValue]
class silkworm.pipelines.JsonLinesPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Append one JSON value per line to a local file.

Parameters:
  • path – Destination JSON Lines file.

  • use_opendal – Use OpenDAL async appends when true and available. None selects OpenDAL automatically when installed.

__init__(path='items.jl', *, use_opendal=None, log_level='DEBUG')[source]
Parameters:
Return type:

None

async open(spider)[source]

Create parent directories and initialize the selected writer.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Flush and close the active local writer.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Serialize and append one item, then return it unchanged.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.LoggedPipeline[source]

Bases: _BatchPipelineMixin

Wrap any item pipeline and control its per-item log level.

Use log_level=None to suppress noisy per-item pipeline messages.

__init__(pipeline, *, log_level='DEBUG')[source]
Parameters:
Return type:

None

property native_batch: bool

Whether the wrapped pipeline opts into engine-managed batching.

async open(spider)[source]

Open the wrapped pipeline.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Close the wrapped pipeline.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Apply the configured log level and delegate item processing.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Delegate native batching when available, otherwise process in order.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.ModelSchema[source]

Bases: Protocol

A Pydantic-style model class: model_validate returns a model instance.

model_validate(obj)[source]

Validate obj and return a model instance, raising on failure.

Parameters:

obj (object)

Return type:

object

__init__(*args, **kwargs)
class silkworm.pipelines.MongoDBPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to a MongoDB collection.

Parameters:
  • connection_string – MongoDB connection URI.

  • database – Database name.

  • collection – Collection receiving documents.

Example:

from silkworm.pipelines import MongoDBPipeline

pipeline = MongoDBPipeline(
    connection_string="mongodb://localhost:27017",
    database="scraping",
    collection="quotes",
)
__init__(connection_string='mongodb://localhost:27017', *, database='scraping', collection='items')[source]

Initialize MongoDBPipeline.

Parameters:
  • connection_string (str) – MongoDB connection string

  • database (str) – Database name

  • collection (str) – Collection name

Return type:

None

async open(spider)[source]

Create the Motor client and select the database and collection.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Close the MongoDB client and release collection references.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Insert a shallow copy so MongoDB cannot add _id to the item.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.MsgPackPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to a file in MessagePack format.

MessagePack is a binary serialization format that is more compact and faster than JSON. This pipeline uses ormsgpack for fast serialization.

Parameters:
  • path – Destination MessagePack file.

  • mode – "write" to replace the file or "append" to retain and extend an existing stream of MessagePack values.

Example:

from silkworm.pipelines import MsgPackPipeline

pipeline = MsgPackPipeline("data/items.msgpack")
# Or append to existing file:
pipeline = MsgPackPipeline("data/items.msgpack", mode="append")

Reading MsgPack files:

import msgpack

# Read all items at once
with open("data/items.msgpack", "rb") as f:
    unpacker = msgpack.Unpacker(f)
    items = list(unpacker)

# Or stream items one by one (memory efficient for large files)
with open("data/items.msgpack", "rb") as f:
    unpacker = msgpack.Unpacker(f)
    for item in unpacker:
        process(item)
__init__(path='items.msgpack', *, mode='write')[source]

Initialize MsgPackPipeline.

Parameters:
  • path (str | Path) – Path to the output file (default: “items.msgpack”)

  • mode (str) – Write mode - “write” (overwrite) or “append” (default: “write”)

Return type:

None

async open(spider)[source]

Open the binary destination in overwrite or append mode.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Flush and close the MessagePack destination.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Serialize and append one MessagePack value.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.MySQLPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to a MySQL database.

Parameters:
  • host – Database host.

  • port – Database port.

  • user – Login user.

  • password – Login password.

  • database – Existing database name.

  • table – Valid unquoted table created automatically when absent.

Example:

from silkworm.pipelines import MySQLPipeline

pipeline = MySQLPipeline(
    host="localhost",
    port=3306,
    user="root",
    password="password",
    database="scraping",
    table="items",
)
__init__(host='localhost', port=3306, user='root', password='', database='scraping', *, table='items')[source]

Initialize MySQLPipeline.

Parameters:
  • host (str) – MySQL host

  • port (int) – MySQL port

  • user (str) – MySQL user

  • password (str) – MySQL password

  • database (str) – Database name

  • table (str) – Table name

Return type:

None

async open(spider)[source]

Create the connection pool and JSON table if absent.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Close the MySQL connection pool and wait for shutdown.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Insert and commit one JSON item with its spider name.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.PolarsPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to a Parquet file using Polars.

Parquet is a columnar storage format optimized for analytics workloads. This pipeline uses Polars for fast and efficient Parquet serialization.

Items are buffered in memory and written when close() runs.

Parameters:
  • path – Destination Parquet file.

  • mode – "write" to replace the file or "append" to merge the buffered items with an existing Parquet file.

Example:

from silkworm.pipelines import PolarsPipeline

pipeline = PolarsPipeline("data/items.parquet")
# Or append to existing file:
pipeline = PolarsPipeline("data/items.parquet", mode="append")

Reading Parquet files:

import polars as pl

# Read entire dataset
df = pl.read_parquet("data/items.parquet")

# Or read with filters/projections (memory efficient)
df = pl.scan_parquet("data/items.parquet").filter(
    pl.col("author") == "John"
).collect()
__init__(path='items.parquet', *, mode='write')[source]

Initialize PolarsPipeline.

Parameters:
  • path (str | Path) – Path to the output file (default: “items.parquet”)

  • mode (str) – Write mode - “write” (overwrite) or “append” (default: “write”)

Return type:

None

async open(spider)[source]

Create parent directories and reset the Parquet item buffer.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Write buffered items, merging an existing file in append mode.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Buffer one item for Parquet serialization during close().

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.PostgreSQLPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to a PostgreSQL database.

Parameters:
  • host – Database host.

  • port – Database port.

  • user – Login user.

  • password – Login password.

  • database – Existing database name.

  • table – Valid unquoted table created automatically when absent.

Example:

from silkworm.pipelines import PostgreSQLPipeline

pipeline = PostgreSQLPipeline(
    host="localhost",
    port=5432,
    user="postgres",
    password="password",
    database="scraping",
    table="items",
)
__init__(host='localhost', port=5432, user='postgres', password='', database='scraping', *, table='items')[source]

Initialize PostgreSQLPipeline.

Parameters:
  • host (str) – PostgreSQL host

  • port (int) – PostgreSQL port

  • user (str) – PostgreSQL user

  • password (str) – PostgreSQL password

  • database (str) – Database name

  • table (str) – Table name

Return type:

None

async open(spider)[source]

Create the async pool and JSONB table if absent.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Close the PostgreSQL connection pool.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Insert one JSONB item with its spider name.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.RssPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to an RSS 2.0 feed (buffered).

Items must provide title, link, and description fields (configurable).

Parameters:
  • path – Destination RSS XML file.

  • channel_title – Feed title.

  • channel_link – Canonical feed or site URL.

  • channel_description – Feed description.

  • max_items – Maximum most-recent valid items retained, or None for no limit.

  • item_title_field – Mapping key containing each item title.

  • item_link_field – Mapping key containing each item URL.

  • item_description_field – Mapping key containing each item description.

  • item_pub_date_field – Optional publication-date key.

  • item_guid_field – Optional GUID key.

  • item_author_field – Optional author key.

__init__(path='feed.xml', *, channel_title, channel_link, channel_description, max_items=50, item_title_field='title', item_link_field='link', item_description_field='description', item_pub_date_field=None, item_guid_field=None, item_author_field=None)[source]
Parameters:
  • path (str | Path)

  • channel_title (str)

  • channel_link (str)

  • channel_description (str)

  • max_items (int | None)

  • item_title_field (str)

  • item_link_field (str)

  • item_description_field (str)

  • item_pub_date_field (str | None)

  • item_guid_field (str | None)

  • item_author_field (str | None)

Return type:

None

async open(spider)[source]

Create parent directories and reset the bounded item buffer.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Build and write the RSS 2.0 document from buffered items.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Normalize and buffer a mapping containing required RSS fields.

Non-mapping items and mappings missing required fields are logged and returned without being added to the feed.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.S3JsonLinesPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to S3 in JSON Lines format using async OpenDAL.

Parameters:
  • bucket – Destination bucket.

  • key – Destination object key.

  • region – AWS region.

  • endpoint – Optional S3-compatible service endpoint.

  • access_key_id – Explicit access key or None for provider discovery.

  • secret_access_key – Explicit secret paired with the access key.

Items are buffered and the object is written when the pipeline closes.

Example:

from silkworm.pipelines import S3JsonLinesPipeline

pipeline = S3JsonLinesPipeline(
    bucket="my-bucket",
    key="data/items.jl",
    region="us-east-1",
)
__init__(bucket, key='items.jl', *, region='us-east-1', endpoint=None, access_key_id=None, secret_access_key=None)[source]

Initialize S3JsonLinesPipeline.

Parameters:
  • bucket (str) – S3 bucket name

  • key (str) – S3 object key (path)

  • region (str) – AWS region

  • endpoint (str | None) – Custom S3 endpoint (for S3-compatible services)

  • access_key_id (str | None) – AWS access key ID (uses env vars if not provided)

  • secret_access_key (str | None) – AWS secret access key (uses env vars if not provided)

Return type:

None

async open(spider)[source]

Create the asynchronous S3 operator and reset the item buffer.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Write buffered JSON Lines to the configured S3 object.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Serialize and buffer one JSON line for upload during close().

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.SFTPPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to an SFTP server in JSON Lines format.

Items are buffered until close(), when a single upload is performed.

Parameters:
  • host – SFTP server hostname.

  • user – Account name used for authentication.

  • password – Password authentication secret. Either this or private_key must be provided.

  • remote_path – Destination path on the SFTP server.

  • port – SFTP server port.

  • private_key – Path to a private key used for authentication.

  • known_hosts – Alternate OpenSSH known-hosts file. By default AsyncSSH uses the current user’s standard known-hosts file.

  • verify_host_key – Whether to reject servers whose host key is untrusted.

Example:

from silkworm.pipelines import SFTPPipeline

pipeline = SFTPPipeline(
    host="sftp.example.com",
    user="username",
    password="password",
    remote_path="data/items.jl",
)

The server’s host key is verified against ~/.ssh/known_hosts by default. Pass known_hosts to use another file, or verify_host_key=False to skip verification (not recommended: it allows man-in-the-middle attacks).

__init__(host, user, password=None, remote_path='items.jl', *, port=22, private_key=None, known_hosts=None, verify_host_key=True)[source]

Initialize SFTPPipeline.

Parameters:
  • host (str) – SFTP server hostname

  • user (str) – SFTP username

  • password (str | None) – SFTP password (optional if using private_key)

  • remote_path (str) – Remote file path (default: “items.jl”)

  • port (int) – SFTP port (default: 22)

  • private_key (str | None) – Path to private key file for key-based authentication (optional)

  • known_hosts (str | PathLike[str] | None) – Path to an OpenSSH known_hosts file used to verify the server’s host key. None (default) uses asyncssh’s default, the user’s ~/.ssh/known_hosts.

  • verify_host_key (bool) – Verify the server’s host key (default: True). Set to False only for trusted networks; it disables protection against man-in-the-middle attacks. Cannot be combined with known_hosts.

Return type:

None

async open(spider)[source]

Reset the in-memory JSON Lines buffer without connecting yet.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Connect, upload all buffered lines, and close SFTP resources.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Serialize and buffer one JSON line for upload during close().

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.SQLitePipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Stream items into a SQLite table as JSON documents.

Parameters:
  • path – SQLite database path.

  • table – Valid unquoted table name created automatically when absent.

__init__(path='items.db', table='items')[source]
Parameters:
Return type:

None

async open(spider)[source]

Open the database and create the destination table if needed.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Commit outstanding writes and close the database connection.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Insert one JSON-serialized item and return it unchanged.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.SnowflakePipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to a Snowflake data warehouse.

Parameters:
  • account – Snowflake account identifier.

  • user – Login user.

  • password – Login password.

  • database – Database name.

  • schema – Schema name.

  • warehouse – Compute warehouse name.

  • table – Valid unquoted destination table.

  • role – Optional active role.

Example:

from silkworm.pipelines import SnowflakePipeline

pipeline = SnowflakePipeline(
    account="myaccount",
    user="myuser",
    password="mypassword",
    database="mydatabase",
    schema="myschema",
    warehouse="mywarehouse",
    table="items",
)
__init__(account, user, password, database, schema, warehouse, *, table='items', role=None)[source]

Initialize SnowflakePipeline.

Parameters:
  • account (str) – Snowflake account identifier

  • user (str) – Snowflake username

  • password (str) – Snowflake password

  • database (str) – Database name

  • schema (str) – Schema name

  • warehouse (str) – Warehouse name

  • table (str) – Table name (default: “items”)

  • role (str | None) – Optional role name

Return type:

None

async open(spider)[source]

Connect and create the VARIANT-backed table if absent.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Close the Snowflake cursor and connection.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Insert and commit one JSON item with its spider name.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.TaskiqPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends scraped items to a Taskiq broker/queue instead of writing to a file.

This allows you to process items asynchronously with Taskiq workers, enabling distributed processing, retries, and other Taskiq features.

Parameters:
  • broker – Taskiq broker whose lifecycle is managed by the pipeline.

  • task – Decorated Taskiq task to enqueue. When provided, task_name is ignored.

  • task_name – Fully qualified name of a task registered with broker. Required when task is omitted.

Example:

from taskiq import InMemoryBroker
from silkworm.pipelines import TaskiqPipeline

broker = InMemoryBroker()

@broker.task
async def process_item(item):
    # Your item processing logic here
    print(f"Processing: {item}")

pipeline = TaskiqPipeline(broker, task=process_item)
# Or: pipeline = TaskiqPipeline(broker, task_name=".:process_item")
__init__(broker, task=None, task_name=None)[source]

Initialize TaskiqPipeline.

Parameters:
  • broker (_AsyncBroker) – A Taskiq AsyncBroker instance (e.g., InMemoryBroker, RedisBroker)

  • task (_TaskiqTask | None) – A decorated task function (created with @broker.task). If provided, task_name is ignored.

  • task_name (str | None) – Full name of the task registered on the broker (e.g., “.:process_item”). Either task or task_name must be provided.

Return type:

None

async open(spider)[source]

Open the pipeline and start the broker if needed.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Close the pipeline and shutdown the broker.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Send the item to the Taskiq broker for processing.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.ValidationPipeline[source]

Bases: _BatchPipelineMixin

Validate items against a schema and drop (or reject) invalid ones.

schema is either a Pydantic-style model class (anything exposing model_validate, such as pydantic.BaseModel subclasses) or a callable that returns the validated item and raises ValueError/TypeError (Pydantic’s ValidationError is a ValueError) for invalid input. Model results are converted back to JSON-compatible dicts with model_dump(mode="json"), so later pipelines receive normalized data.

Invalid items raise DropItem with reason "invalid" (counted under items_dropped), or propagate the validation error with on_invalid="raise". Combine with the engine’s max_item_drop_rate or min_items to fail crawls when a site redesign breaks extraction. The first log_limit failures are logged with details.

Parameters:
  • schema – Model class or validator callable.

  • on_invalid – "drop" to discard invalid items, "raise" to fail the emitting callback.

  • log_limit – Invalid items logged at warning level before going quiet.

valid

Number of items that passed validation.

invalid

Number of items that failed validation.

__init__(schema, *, on_invalid='drop', log_limit=10)[source]
Parameters:
  • schema (ItemSchema)

  • on_invalid (Literal['drop', 'raise'])

  • log_limit (int)

Return type:

None

async open(spider)[source]

Reset the validation counters.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Log how many items passed and failed validation.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Return the validated item, or drop/raise when it is invalid.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

class silkworm.pipelines.VortexPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to a Vortex file using the vortex-data library.

Vortex is a next-generation columnar file format optimized for high-performance data processing with 100x faster random access reads compared to Parquet, 10-20x faster scans, and similar compression ratios. It provides zero-copy compatibility with Apache Arrow.

Items are buffered in memory and converted to an Arrow table during close().

Parameters:

path – Destination Vortex file.

Example:

from silkworm.pipelines import VortexPipeline

pipeline = VortexPipeline("data/items.vortex")

Reading Vortex files:

import vortex

# Open and read a Vortex file
vortex_file = vortex.file.open("data/items.vortex")
arrow_table = vortex_file.to_arrow().read_all()

# Or convert to other formats
df = vortex_file.to_polars()  # Polars DataFrame
df = vortex_file.to_dataset().to_table()  # PyArrow Table
__init__(path='items.vortex')[source]

Initialize VortexPipeline.

Parameters:

path (str | Path) – Path to the output file (default: “items.vortex”)

Return type:

None

async open(spider)[source]

Create parent directories and reset the Vortex item buffer.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Convert buffered mappings to Arrow and write the Vortex file.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Buffer one item for Vortex serialization during close().

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.WebhookPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that sends items to a webhook endpoint using the wreq HTTP client.

This pipeline uses the same HTTP client (wreq) as the spider itself for sending data to webhooks, ensuring consistent behavior and browser impersonation.

Parameters:
  • url – Webhook endpoint URL.

  • method – HTTP method used for deliveries.

  • headers – Headers included with every delivery.

  • timeout – Per-delivery timeout, either in seconds or as a duration.

  • batch_size – Number of items per request. A value of 1 sends items immediately; a partial final batch is sent by close().

Example:

from silkworm.pipelines import WebhookPipeline

pipeline = WebhookPipeline(
    url="https://webhook.site/unique-id",
    method="POST",
    headers={"Authorization": "Bearer token123"},
)
__init__(url, *, method='POST', headers=None, timeout=30.0, batch_size=1)[source]

Initialize WebhookPipeline.

Parameters:
  • url (str) – Webhook endpoint URL

  • method (str) – HTTP method (default: “POST”)

  • headers (dict[str, str] | None) – Optional HTTP headers to send with each request

  • timeout (float | timedelta | None) – Request timeout in seconds (default: 30.0)

  • batch_size (int) – Number of items to batch before sending (default: 1 for immediate sending)

Return type:

None

async open(spider)[source]

Create the webhook client and reset the item batch.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Send a partial batch and close the webhook client.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Buffer one item and send when batch_size is reached.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.XMLPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Stream mapping items into one XML document.

Parameters:
  • path – Destination XML file.

  • root_element – Document root tag.

  • item_element – Tag wrapping each item.

__init__(path='items.xml', *, root_element='items', item_element='item')[source]
Parameters:
Return type:

None

async open(spider)[source]

Open the destination and write the XML declaration and root tag.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Write the closing root tag and close the XML file.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Append one mapping as an XML element and return it unchanged.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

class silkworm.pipelines.YAMLPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that writes items to a YAML file.

Parameters:

path – Output YAML path. Items are buffered and written as one sequence when the pipeline closes.

Example:

from silkworm.pipelines import YAMLPipeline

pipeline = YAMLPipeline("data/items.yaml")
__init__(path='items.yaml')[source]

Initialize YAMLPipeline.

Parameters:

path (str | Path) – Path to the output file (default: “items.yaml”)

Return type:

None

async open(spider)[source]

Create parent directories and reset the YAML item buffer.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Write all buffered items as a YAML sequence.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Buffer one item for YAML serialization during close().

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]

type silkworm.pipelines.ZenohKeyResolver = Callable[[JSONValue, Spider], str | Awaitable[str]]
class silkworm.pipelines.ZenohPipeline[source]

Bases: _BatchPipelineMixin

native_batch = True

Pipeline that publishes JSON-serialized items to Zenoh.

Parameters:
  • key_expr – Static Zenoh key expression or a sync/async resolver receiving (item, spider) and returning one.

  • config – Configuration used to open a pipeline-owned session. Mutually exclusive with session; a default configuration is used when both are omitted.

  • session – Existing session to reuse. The pipeline never closes an injected session.

  • encoding – Encoding declared for published payloads.

  • congestion_control – Behavior when Zenoh’s transmission queue is full.

  • priority – Transmission priority for published items.

  • express – Whether messages should bypass batching where possible.

  • reliability – Reliability requested from declared publishers.

  • allowed_destination – Locality constraint for published items.

Publishers created by this pipeline are cached by key expression and undeclared during close(). Zenoh’s synchronous calls run in worker threads so they do not block Silkworm’s event loop.

Example:

from silkworm.pipelines import ZenohPipeline

pipeline = ZenohPipeline("scraping/items")
__init__(key_expr, *, config=None, session=None, encoding='application/json', congestion_control=None, priority=None, express=None, reliability=None, allowed_destination=None)[source]

Initialize a Zenoh publisher pipeline.

Parameters:
  • key_expr (str | ZenohKeyResolver)

  • config (Config | None)

  • session (Session | None)

  • encoding (str | Encoding)

  • congestion_control (CongestionControl | None)

  • priority (Priority | None)

  • express (bool | None)

  • reliability (Reliability | None)

  • allowed_destination (Locality | None)

Return type:

None

async open(spider)[source]

Open or attach to a session and declare a static publisher.

Parameters:

spider (Spider)

Return type:

None

async close(spider)[source]

Undeclare publishers and close a pipeline-owned session.

Parameters:

spider (Spider)

Return type:

None

async process_item(item, spider)[source]

Serialize and publish one item, then pass it to the next pipeline.

Parameters:
  • item (JSONValue)

  • spider (Spider)

Return type:

JSONValue

async process_items(items, spider)[source]

Process a batch sequentially, preserving single-item semantics.

Parameters:
Return type:

list[JSONValue]