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)
- 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
- type silkworm.pipelines.BatchItemCallback = Callable[[list[JSONValue], Spider], list[JSONValue] | Awaitable[list[JSONValue] | None] | None]¶
- class silkworm.pipelines.BatchItemPipeline[source]¶
Bases:
ItemPipeline,ProtocolItem pipeline that can process an explicit batch in one call.
- 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.
- 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
- class silkworm.pipelines.CallbackPipeline[source]¶
Bases:
_BatchPipelineMixinPipeline 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). ReturningNonepreserves 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
- 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
- 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
- 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", )
- 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
- 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
idhash key.region_name – AWS region.
aws_access_key_id – Explicit access key, or
Nonefor 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
- 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
AsyncElasticsearchclient 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.
- 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
- 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")
- 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
- 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", )
- 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
- 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:
- 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
- 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
streamduring startup when needed.create_topic_if_not_exists – Create
topicduring startup when needed.topic_partitions_count – Number of partitions used when creating
topic.send_retries – Number of producer send retries, or
Nonefor 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
- type silkworm.pipelines.ItemCallback = Callable[[JSONValue, Spider], JSONValue | Awaitable[JSONValue | None] | None]¶
- class silkworm.pipelines.ItemPipeline[source]¶
Bases:
ProtocolProtocol 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¶
- 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.
Noneselects OpenDAL automatically when installed.
- 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
- class silkworm.pipelines.LoggedPipeline[source]¶
Bases:
_BatchPipelineMixinWrap any item pipeline and control its per-item log level.
Use
log_level=Noneto suppress noisy per-item pipeline messages.- __init__(pipeline, *, log_level='DEBUG')[source]¶
- Parameters:
pipeline (ItemPipeline)
log_level (LogLevel)
- Return type:
None
- 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
- class silkworm.pipelines.ModelSchema[source]¶
Bases:
ProtocolA Pydantic-style model class:
model_validatereturns a model instance.- __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.
- 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
- 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)
- 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
- 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.
- 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
- 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()
- 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
- 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.
- 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
- 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
Nonefor 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:
- 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
- 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
Nonefor 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
- 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_keymust 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_hostsby default. Passknown_hoststo use another file, orverify_host_key=Falseto 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
- 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.
- 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
- 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:
- 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
- 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_nameis ignored.task_name – Fully qualified name of a task registered with
broker. Required whentaskis 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
- class silkworm.pipelines.ValidationPipeline[source]¶
Bases:
_BatchPipelineMixinValidate items against a schema and drop (or reject) invalid ones.
schemais either a Pydantic-style model class (anything exposingmodel_validate, such aspydantic.BaseModelsubclasses) or a callable that returns the validated item and raisesValueError/TypeError(Pydantic’sValidationErroris aValueError) for invalid input. Model results are converted back to JSON-compatible dicts withmodel_dump(mode="json"), so later pipelines receive normalized data.Invalid items raise
DropItemwith reason"invalid"(counted underitems_dropped), or propagate the validation error withon_invalid="raise". Combine with the engine’smax_item_drop_rateormin_itemsto fail crawls when a site redesign breaks extraction. The firstlog_limitfailures 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.
- async open(spider)[source]¶
Reset the validation counters.
- Parameters:
spider (Spider)
- Return type:
None
- 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
- 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
- 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
1sends items immediately; a partial final batch is sent byclose().
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
- 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.
- 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
- 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")
- 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
- 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