Create your first pipeline
Plombery discovers your pipelines from a pipelines/ folder and runs them with
the plombery command. A minimal project is just that folder with one file in
it:
Glossary
Before starting, let's define some naming so there will be no confusion!
- Task: a python function that performs some job, it's the base block for building a pipeline
- Pipeline: a graph of one or more Tasks, a pipeline can be run via a schedule, manually, etc.
- Trigger: is the entrypoint to run a pipeline, a trigger can be a schedule, a webhook, a button on the web UI, etc.
- Pipeline Run: (sometimes simply referred as Run) is the result of running a pipeline
Write a pipeline
Create pipelines/sales.py. A Task is a Python function decorated with
@task, and a Pipeline groups the tasks defined inside its with block:
from datetime import datetime
from random import randint
from plombery import Pipeline, task, get_logger
with Pipeline(id="sales_pipeline") as pipeline:
@task
async def fetch_raw_sales_data():
"""Fetch latest 50 sales of the day"""
# Using Plombery's logger, your logs are stored and shown in the web UI
logger = get_logger()
logger.debug("Fetching sales data...")
sales = [
{
"price": randint(1, 1000),
"store_id": randint(1, 10),
"date": datetime.today(),
"sku": randint(1, 50),
}
for _ in range(50)
]
logger.info("Fetched %s sales data rows", len(sales))
# Returning a value stores it and makes it available in the web UI;
# it's also passed to any task downstream of this one
return sales
Every task defined inside the with Pipeline() block is added to the pipeline
automatically, and the pipeline registers itself with Plombery when the block
ends — importing the file is all it takes. See Pipelines for
how to add more tasks and declare dependencies between them with >>.
Run it
plombery run imports every file in the pipelines/ folder, so each
register_pipeline runs, and serves the web app. Open
http://localhost:8000 and you'll find
your pipeline, ready to run manually.
By default it looks for a pipelines/ folder in the current directory and
binds to 127.0.0.1:8000. Change any of that:
While you're writing pipelines, --reload restarts Plombery
whenever a file changes, so a new or edited pipeline shows up without you
stopping the server. A file that doesn't import yet is reported and skipped,
so a half-written pipeline doesn't take the app down with it:
Only the folder you run from is watched. If your pipelines import a library
that lives elsewhere, add it with --reload-dir path/to/it, as many times as
you need.
Without watchfiles installed, reloading falls back to polling: edits are
still picked up promptly, but a pipeline file you add only shows up once you
edit a file that already existed, and polling slows down if a virtualenv sits
in a watched folder. Install it to avoid both:
Reloading is for developing pipelines only — don't use it to serve Plombery in production, where a file that fails to import should stop the deploy rather than be skipped.
Schedule it
A pipeline with no trigger can be run manually from the web UI or its HTTP
endpoint. To have Plombery run it on a schedule, add a Trigger:
from apscheduler.triggers.interval import IntervalTrigger
from plombery import Pipeline, Trigger, task
with Pipeline(
id="sales_pipeline",
description="Aggregate sales activity from all stores across the country",
triggers=[
Trigger(
id="daily",
name="Daily",
description="Run the pipeline every day",
schedule=IntervalTrigger(days=1),
),
],
) as pipeline:
@task
async def fetch_raw_sales_data():
...
See Triggers for cron schedules and triggers with parameters.
Without the CLI
plombery run is the quickest way to start, but Plombery is a FastAPI app, so
you can also run it yourself — useful to embed it in a larger app, or to use a
different ASGI server. Import your pipelines, expose the app with get_app(),
and serve it: