Orchestration
The pipeline canvas, step reference, schedules, triggers, variables and secrets.
A pipeline is one graph of steps on a canvas. Some steps move data, and the rest run a query, call an API, branch on a condition, and so on. Lines between steps control the order they run in. When you deploy a pipeline, it's compiled into a workflow that runs on Pipeloom's execution engine.
Moving data: sources, destinations and syncs
Data moves through three pieces on the canvas:
- A Source step reads from a source connector in your workspace.
- A Destination step writes to a destination connector.
- The line from a source to a destination is the sync. It copies the streams (tables) you choose from one to the other. Select the line to pick an existing sync between those two connectors, or set up a new one and choose its streams.
Steps connected after a destination start once the data has landed there. Several sources can feed one destination; each line is its own sync.
Connectors are shared across the workspace, so one database needs only one connector, used by as many pipelines as you like. Syncs are different.
A sync belongs to one pipeline
Each sync runs in exactly one pipeline, and that pipeline's schedule decides when it runs. If a sync already runs in another pipeline, you can't put it on this canvas too. The canvas tells you where it runs and offers to react to it instead.
To act on data that another pipeline's sync brings in, add a Sync from another pipeline step and pick that sync. It moves nothing itself: each time that sync lands new data, this pipeline starts and runs the steps after it. The sync keeps running in its own pipeline, on that pipeline's schedule.
Don't set up a second sync from the same source just to get a copy of its data. On a change-data-capture (CDC) source, both syncs would read from the same replication slot and take each other's changes. Pipeloom refuses a new sync from a source whose slot is already in use, and points you to the sync that already reads it.
Detaching a sync from a pipeline keeps the sync and its history; the pipeline just stops running it.
Step reference
Every step below is in the Add component palette.
Connectors
| Step | What it does |
|---|---|
Source (conduit.source) | A source connector from your workspace. Connect it to a destination to sync its data there. |
Destination (conduit.destination) | A destination connector from your workspace. Every source connected to it syncs into it. Has a timeout (default 180 minutes) after which the step fails if a sync into it hasn't finished. |
Sync from another pipeline (conduit.listen) | Reacts to a sync that another pipeline runs: starts this pipeline each time that sync lands new data. |
Transform
| Step | What it does |
|---|---|
sql.query | Runs SQL against a warehouse, using a credential stored in Pipeloom. |
dbt.build | Runs dbt build against a project directory, with an optional selector and full refresh. |
Compute & integration
| Step | What it does |
|---|---|
script.shell | Runs shell commands. |
http.request | Calls an HTTP endpoint: a webhook, a reverse-ETL trigger, an alert. |
notify.slack | Posts a message to Slack through an incoming webhook. |
pipeline.run | Runs another pipeline as a step, optionally waiting for it to finish. |
component.custom | Runs a custom Python component from your workspace. |
Flow control
| Step | What it does |
|---|---|
flow.start | Marks an entry point on the diagram. Any step with no incoming line runs first anyway; this just labels it. |
flow.if | Branches on a condition: a "then" path when it's true, and an optional "else" path when it's false. |
flow.foreach | Runs a group of steps once per item in a list, in parallel, one at a time, or up to a limit you set. |
flow.wait | Pauses this path for a fixed duration (for example PT5M), then continues. |
flow.and | Waits for every incoming line to succeed before continuing. A step with several incoming lines already does this; "And" makes it visible. |
flow.or | Joins several incoming lines. Continues once, as soon as the first branch succeeds, or waits for every branch, as you choose. |
flow.end.success | Marks a path as deliberately finished. |
flow.end.failure | Fails the pipeline run when this path is reached, with a message you choose. |
A step's outputs can be referenced from later steps, so a downstream sql.query or
component.custom can use what an earlier step produced.
Pipelines built before sources and destinations were steps on the canvas may contain an older
Pipeloom sync (conduit.sync) step. These still open and run, but the palette no longer
offers it. Use a source, a destination and the line between them instead.
Schedules
A schedule runs the whole pipeline, including its syncs; nothing else schedules a sync. Choose a preset (every hour, daily, weekdays, weekly, monthly), write your own cron expression with a timezone, or leave the pipeline to run only when you start it or a trigger fires. You can pause a schedule without deleting it.
Near real-time
The Near real-time preset starts the pipeline every minute:
- If a run is still going when the next minute comes, that minute is skipped. The page shows how often it really runs, based on how long recent runs take.
- Steps after a sync run only when that sync brought new rows. Otherwise they show as skipped.
- Every stream in the pipeline's syncs must be incremental. Deploying checks this and refuses a stream set to full refresh, since it would be copied in full every minute.
- A source that reads changes by cursor rather than CDC works, but you'll be warned: every run rescans its tables, and deleted rows are never picked up. Switch the source to CDC for true near-real-time.
Triggers
A trigger starts the pipeline when a sync finishes, instead of guessing when the data will have landed. Pick the sync and when to fire: when it succeeds, when it fails, or when it landed new data.
Triggers are for syncs that no pipeline runs. For a sync that runs in another pipeline, use a Sync from another pipeline step instead. A pipeline also can't trigger on a sync it runs itself, since every run would start another.
Triggers need a one-time setup by an administrator for each workspace. Until then, the trigger panel says they won't fire.
Notifications
A pipeline can post to a webhook when a run starts, succeeds or fails. Notifications are set on the pipeline, not per step. The webhook URL is kept as a secret.
Variables and secrets
Two scopes of configuration are available to steps at runtime:
- Workspace variables: defaults for the whole workspace (for example, a target schema name or a log level), available to every pipeline in it.
- Pipeline variables: values that override a workspace variable with the same key, for that one pipeline.
Either can be marked secret. A secret's value is never stored with your pipeline definition or shown again in the UI once saved. It's resolved only when the pipeline runs, the same way workspace secrets are.
Variables are passed as environment variables to every step that runs code: script.shell,
sql.query, dbt.build, and component.custom. Steps that don't run your code (sources,
destinations, http.request, notify.slack, flow.*) don't use them.
Building components that use variables and secrets? The Custom Components page shows how to read them from Python, and describes a current gap to know about before you rely on it.