Last updated

Pipelines and Automation: Building Smart Workflows 🔄

Learn how to create and manage pipelines in Dataloop - your key to automating workflows and data processing.

Dataloop Login 🔐

import dtlpy as dl
from dtlpy.entities.node import PipelineNodeIO

# Interactive login — opens a browser window
if dl.token_expired():
    dl.login()

Project Setup

# Set your project and dataset names
project_name = "onboarding-project"

try:
    # Try to get existing project
    project = dl.projects.get(project_name=project_name)
    print(f"Project '{project_name}' already exists")
except dl.exceptions.NotFound:
    project = dl.projects.create(project_name=project_name)
    # Create project if it doesn't exist
    print(f"Created project '{project_name}'")

Creating Code Node Pipeline 🚀

Pipeline Flow Diagram

┌─────────────────┐       ┌─────────────────┐       ┌─────────────────┐
│   DatasetNode   │       │    CodeNode     │       │   DatasetNode   │
│ (source-dataset)│ ----> │ (process-item)  │ ----> │(output-dataset) │
|                 |       |                 |       |                 |
└─────────────────┘       └─────────────────┘       └─────────────────┘

Code Node Processing Function

# Define processing method for CodeNode
def process_item(item: dl.Item):
    """
    Process an item - add metadata and return it.
    This method runs directly in the pipeline, no service deployment needed.
    """

    import datetime

    # Initialize metadata
    item.metadata = item.metadata or {}
    item.metadata['user'] = item.metadata.get('user', {})

    # Add processing metadata
    item.metadata['user']['processed'] = True
    item.metadata['user']['timestamp'] = datetime.datetime.now().isoformat()
    item.metadata['user']['pipeline'] = 'code-node-pipeline'

    # Update item in Dataloop
    item.update(system_metadata=True)

    print(f'Processed item: {item.name}')
    return item

Creating Pipelines

# Create pipeline
pipeline = project.pipelines.create(name='code-node-pipeline')
print(f'Pipeline created: {pipeline.name} ({pipeline.id})')

Getting Existing Pipelines

# Get pipeline by ID
pipeline = project.pipelines.get(pipeline_id=pipeline.id)

# List all pipelines in project
project.pipelines.list()

Pipeline Constructions

# Create datasets
source_dataset = project.datasets.create(dataset_name='code-node-source')
output_dataset = project.datasets.create(dataset_name='code-node-output')

# Create nodes

# 1. Source dataset node
source_node = dl.DatasetNode(
    name='source-dataset',
    project_id=project.id,
    dataset_id=source_dataset.id,
    position=(1, 1)
)

# 2. CodeNode - runs the Python method directly
code_node = dl.CodeNode(
    name='process-item',
    position=(2, 1),
    project_id=project.id,
    method=process_item,
    project_name=project.name
)

# 3. Output dataset node
output_node = dl.DatasetNode(
    name='output-dataset',
    project_id=project.id,
    dataset_id=output_dataset.id,
    position=(3, 1)
)

print('Nodes created successfully')

Installing the Pipeline

# Connect nodes: Source → CodeNode → Output
pipeline.nodes.add(node=source_node).connect(node=code_node).connect(node=output_node)

pipeline.update()
pipeline.install()
print(f'Pipeline ready: {pipeline.name} ({pipeline.id})')

# Open pipeline in Dataloop platform
pipeline.open_in_web()

Congratulations! Your First Pipeline is Ready 🎉

Testing the Pipeline

# Upload a test item to source dataset
# Replace with your actual file path
test_item = source_dataset.items.upload(
    local_path=r"/path/to/your/image.jpg",
    remote_path='/'
)
print(f'Uploaded test item: {test_item.name}')
# Execute pipeline manually (optional)
pipeline_execution = pipeline.pipeline_executions.create(
    pipeline_id=pipeline.id,
    execution_input=[dl.FunctionIO(type=dl.PackageInputType.ITEM, value=test_item.id, name='item')]
)
print(f'Execution started: {pipeline_execution.id}')

pipeline.open_in_web()

Creating a Model Pipeline

Pipeline Flow Diagram

┌─────────────────┐       ┌─────────────────────┐
│   DatasetNode   │       │  FunctionNode (ML)  │
│ (source-dataset)│ ----> │ (mobilenet-predict) │
│  Position: (1,1)│       │   Position: (2,1)   │
└─────────────────┘       └─────────────────────┘

Creating a Pipeline

# Create a new pipeline
pipeline = project.pipelines.create(name='model-pipeline')

# Print pipeline details
print(pipeline)

Pipeline Ingredients

# Create datasets
dataset_source = project.datasets.create(dataset_name='ds-model-source')
print(f'Created dataset: {dataset_source.name}')
# Check if model is already installed on project
model_name = "mobilenet"
try:
    model = project.models.get(model_name=model_name)
except dl.exceptions.NotFound:
    model_dpk = dl.dpks.get(dpk_name="mobilenet")
    print(f"App '{model_dpk.name}' not found, installing...")
    model_app = project.apps.install(dpk=model_dpk)
    model = project.models.get(model_name=model_name)
# Deploy model with default configuration
if not model.services.list().items:
   print("No services found for the model")
   model_service = model.deploy()
else:
   print("Services found for the model")
   model_service = model.services.list().items[0]
   print(f"Service found: {model_service.name}")

service_id = model_service.id

# Print service details
model_service.print()
print(f'Service deployed: {model_service.name}')

Pipeline Constructions

# Dataset source
dataset_node = dl.DatasetNode(
    name='source-dataset',
    project_id=project.id,
    dataset_id=dataset_source.id,
    position=(1, 1),
)

# Create ML predict node
# - Set node_type to ML
# - Link to the model via metadata.modelId
# - Link to the installed app for UI function picker resolution
# - Use standard Item -> Item ports
# - function_name is a model lifecycle action: train | predict | evaluate | embed
app = next(a for a in project.apps.list().all() if a.dpk_name == 'mobilenet')
item_port = PipelineNodeIO(
    input_type=dl.PackageInputType.ITEM, name='item', display_name='item', actions=[],
)

predict_node = dl.FunctionNode(
    name='mobilenet-predict',
    service=model_service,
    function_name='predict',
    project_id=project.id,
    position=(2, 1),
)
predict_node.node_type = dl.PipelineNodeType.ML
predict_node.metadata['modelId'] = model.id
predict_node.app_id, predict_node.app_name, predict_node.dpk_name = app.id, app.name, app.dpk_name
predict_node.inputs = [item_port]
predict_node.outputs = [item_port]

print(f'Predict node ready: model={model.name!r}  app={app.name!r}')

Installing the Pipeline

# Connect nodes and install pipeline
pipeline.nodes.add(node=dataset_node).connect(node=predict_node)
pipeline.update()
pipeline.install()
print(f'Pipeline ready: {pipeline.name} ({pipeline.id})')

Testing the Pipeline

item = dataset_source.items.upload(
    local_path=r"/path/to/your/image.jpg",
    remote_path='/'
)

print(f"Uploaded: {item.name} ({item.id})")

pipeline_execution = pipeline.pipeline_executions.create(
    pipeline_id=pipeline.id,
    execution_input=[dl.FunctionIO(type=dl.PackageInputType.ITEM, value=item.id, name='item')]
)
print(f"Execution started: {pipeline_execution.id}")

pipeline.open_in_web()

Pipeline Management 📋

1. Basic Operations

# These are reference examples — replace with your actual pipeline
# project.pipelines.delete(pipeline_id=pipeline.id)
# project.pipelines.pause(pipeline=pipeline)
# project.pipelines.reset(pipeline=pipeline)

# Open pipeline in web UI
pipeline.open_in_web()

2. Pipeline Monitoring

# Get pipeline statistics
project.pipelines.stats(pipeline='pipeline_entity')

# Get pipeline execution object
pipeline_executions = pipeline.pipeline_executions.get(pipeline_id='pipeline_id')

# List project pipeline executions
pipeline.pipeline_executions.list()

Ready to bring it all together? Let's move on to the final chapter! 🚀