How to Use Conditional Task Execution with the Condition Parameter in Sieves
You can skip documents on a per-task basis by passing a callable to the condition parameter of any Task in Sieves, which evaluates each document and stores None for skipped items while preserving the original document order.
Sieves is a modular document processing framework that supports conditional task execution through a flexible condition parameter available on all tasks. This feature allows you to filter documents dynamically based on custom logic, ensuring that expensive operations like model inference only run on relevant documents. By leveraging the condition parameter in the mantisai/sieves repository, you can build efficient pipelines that skip unnecessary processing without breaking the document stream.
How Conditional Task Execution Works in Sieves
The Task Base Class and Condition Storage
The conditional logic is implemented in the base Task class defined in sieves/tasks/core.py. When initializing any task, the condition parameter is stored as self._condition:
# sieves/tasks/core.py
class Task(abc.ABC):
def __init__(self, task_id, include_meta, batch_size, condition=None):
# ...
self._condition = condition # ← stored here
This design ensures that all subclasses—including built-in tasks like Classification and Chunking—inherit conditional execution capabilities without additional implementation.
The Execution Flow and Document Filtering
When a pipeline invokes a task, the __call__ method (lines 68-88 in sieves/tasks/core.py) handles conditional filtering automatically. The method evaluates the condition for each document in the batch to determine which items to process:
# sieves/tasks/core.py (lines 68-88)
while docs_batch := [doc for doc in itertools.islice(docs, batch_size)]:
passing_indices = {
idx for idx, doc in enumerate(docs_batch)
if self._condition is None or self._condition(doc)
}
processed = self._call(d for i, d in enumerate(docs_batch) if i in passing_indices)
for idx, doc in enumerate(docs_batch):
if idx in passing_indices:
yield next(processed) # run the task
else:
doc.results[self.id] = None # skipped → store None
yield doc
Documents that fail the condition check receive None in their results dictionary under the task ID, while passing documents are processed normally. The original document order is preserved throughout the pipeline.
The pipeline itself simply forwards documents to each task. Because each task handles its own conditional logic in sieves/pipeline/core.py, no extra pipeline configuration is required:
# sieves/pipeline/core.py
for i, task in enumerate(self._tasks):
processed_docs = task(processed_docs) # ← each task decides what to do
When to Use Conditional Task Execution
Conditional task execution is particularly valuable when processing heterogeneous document collections. Common scenarios include:
| Situation | Example Condition |
|---|---|
| Skip short texts | lambda d: len(d.text or "") > 20 |
| Process only PDFs | lambda d: d.uri and d.uri.endswith('.pdf') |
| Rate limiting | lambda d, i=itertools.count(): next(i) < N |
| External data validation | Any callable that queries external APIs or databases to determine eligibility |
Implementation Examples
Custom Task with a Condition
The following example demonstrates a custom task that only processes documents longer than 20 characters. This pattern is tested in sieves/tests/tasks/test_conditional_execution.py.
from sieves.tasks.core import Task
from sieves.data import Doc
class DummyTask(Task):
def __init__(self, condition=None):
super().__init__(task_id="DummyTask", include_meta=False,
batch_size=-1, condition=condition)
def _call(self, docs):
for doc in docs:
doc.results[self.id] = {"processed": True}
yield doc
To use this task with conditional execution:
from sieves.pipeline import Pipeline
from sieves import Doc
docs = [Doc(text="short"), Doc(text="this is a much longer document")]
task = DummyTask(condition=lambda d: len(d.text or "") > 20)
pipe = Pipeline([task])
for d in pipe(docs):
print(d.results["DummyTask"])
# Output: None for the short doc, {"processed": True} for the long one
Built-in Classification with Conditional Logic
Built-in predictive tasks like Classification inherit the condition parameter from the base Task class. This example from sieves/tests/docs/test_pipeline_docs.py shows how to conditionally classify only long documents.
from sieves import Pipeline, tasks, Doc
from sieves.models import SmallTransformer # any ModelWrapper implementation
model = SmallTransformer()
docs = [
Doc(text="short"),
Doc(text="this is a much longer document that will be processed"),
Doc(text="med")
]
def is_long(doc: Doc) -> bool:
return len(doc.text or "") > 20
# Pass condition directly to the task constructor
task = tasks.Classification(labels=["science", "politics"],
model=model,
condition=is_long)
pipe = Pipeline([task])
for d in pipe(docs):
print(d.results[task.id]) # None for short/med docs, prediction dict for the long one
Multiple Tasks with Independent Conditions
You can assign different conditions to different tasks within the same pipeline. Each task evaluates its condition independently, allowing complex processing workflows.
from sieves import Pipeline, tasks, Doc
from sieves.models import SmallTransformer
from sieves.chunkers import example_chunker
model = SmallTransformer()
docs = [
Doc(text="short"),
Doc(text="this is a much longer document"),
Doc(text="medium text here")
]
# Task 1: Only chunk documents longer than 10 characters
task1 = tasks.Chunking(example_chunker,
condition=lambda d: len(d.text or "") > 10)
# Task 2: Only classify documents longer than 20 characters
task2 = tasks.Classification(labels=["science", "politics"],
model=model,
condition=lambda d: len(d.text or "") > 20)
pipe = Pipeline([task1, task2])
for d in pipe(docs):
print(d.chunks, d.results.get(task2.id))
Key Implementation Files
The conditional execution system is implemented across these core files in the mantisai/sieves repository:
| File | Purpose | Link |
|---|---|---|
sieves/tasks/core.py |
Defines the base Task class with condition parameter and the __call__ method that filters documents per batch. |
View source |
sieves/pipeline/core.py |
Orchestrates task execution; forwards documents to each task which handles its own conditional logic independently. | View source |
sieves/tests/tasks/test_conditional_execution.py |
Unit tests demonstrating custom task behavior with conditional execution. | View source |
sieves/tests/docs/test_pipeline_docs.py |
Integration examples showing built-in tasks with conditions in pipeline contexts. | View source |
sieves/tasks/predictive/classification/core.py |
Example of a predictive task inheriting conditional capabilities from the base Task class. |
View source |
Summary
- The
conditionparameter in Sieves accepts any callable that takes aDocobject and returns a boolean, enabling conditional task execution on a per-document basis. - When a document fails the condition check, the task stores
Noneindoc.results[task.id]and yields the document unchanged, preserving the original stream order. - All tasks—including custom subclasses and built-in predictive tasks like
Classification—inherit this functionality from the baseTaskclass insieves/tasks/core.py. - Conditions are evaluated independently for each task in a pipeline, allowing different filtering logic for chunking, classification, or custom processing steps.
- Skipped documents remain in the pipeline for subsequent tasks, making the system cache-friendly and composable.
Frequently Asked Questions
What happens to documents that don't meet the condition?
Documents that fail the condition check are not processed by the task's _call method. Instead, the task automatically sets doc.results[self.id] = None and yields the document unchanged. This ensures the document remains in the pipeline for subsequent tasks while maintaining the original order of the document stream.
Can I use different conditions for different tasks in the same pipeline?
Yes, each task in a Sieves pipeline maintains its own independent condition callable. You can configure Task A to process only PDFs while Task B processes only long documents, and both will execute their conditions separately as the pipeline forwards documents through each step.
Does the condition parameter affect batch processing?
The condition is evaluated within the batch processing loop in Task.__call__. For each batch, the task identifies which documents satisfy the condition (passing_indices) and only processes those, while still yielding skipped documents to maintain stream continuity. This means batching efficiency is preserved while allowing per-document filtering.
How do I access the document index in my condition function?
The condition callable receives only the Doc object as its argument by default. If you need access to the document index, you can use a closure with itertools.count() or a class-based callable that maintains state. For example: lambda d, i=itertools.count(): next(i) < 100 processes only the first 100 documents.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →