Understanding these concepts will help you build effective data pipelines.
A pipeline is a complete data flow definition. It contains:
{
version: 1,
steps: [...],
edges: [...],
context: { ... },
capabilities: { ... },
hooks: { ... },
}
Steps are the building blocks of a pipeline. Each step has:
| Type | Purpose |
|---|---|
TRIGGER |
Defines how the pipeline starts (MANUAL, SCHEDULE, WEBHOOK, EVENT, FILE, MESSAGE) |
EXTRACT |
Pulls data from external sources |
TRANSFORM |
Modifies, validates, or enriches data |
VALIDATE |
Validates records against rules or schemas |
ENRICH |
Adds data from external lookups |
ROUTE |
Splits data flow based on conditions |
GATE |
Pauses pipeline execution for human approval. Supports MANUAL, THRESHOLD, and TIMEOUT approval types |
LOAD |
Creates or updates Vendure entities |
EXPORT |
Sends data to external destinations |
FEED |
Generates product feeds (Google, Meta, etc.) |
SINK |
Indexes data to search engines |
Edges connect steps to define execution order:
{
from: 'extract-step',
to: 'transform-step',
branch: 'optional-branch-name' // For routing
}
Data flows from the from step to the to step. Multiple edges can originate from one step (fan-out) and multiple edges can target one step (fan-in).
A record is a single data item flowing through the pipeline. Records are typically JSON objects:
{
"id": "123",
"name": "Product Name",
"price": 29.99,
"category": "Electronics"
}
Records are processed individually through transform and load steps. Extractors typically yield multiple records from a single source.
Adapters are reusable components that perform specific operations:
| Adapter Type | Purpose |
|---|---|
| Extractor | Connects to data sources and yields records |
| Operator | Transforms individual record fields |
| Loader | Creates or updates Vendure entities |
| Exporter | Writes data to external destinations |
| Feed Generator | Creates product feed files |
| Search Sink | Indexes data to search engines |
Connections store reusable configuration for external systems:
Connections are referenced by code in extract and export steps:
.extract('fetch-api', {
adapterCode: 'httpApi',
connectionCode: 'my-api', // References saved connection
url: '/products',
})
Secrets store sensitive values like API keys and passwords. They are:
.extract('api-call', {
adapterCode: 'httpApi',
connectionCode: 'catalog-api',
url: '/products',
auth: {
type: 'BEARER',
secretCode: 'api-key', // References saved secret
},
})
Secret-backed HTTP authentication must reference a saved connection with a base URL. Relative paths resolve against that URL; absolute URLs and redirects must keep the connection’s exact origin.
A run is a single execution of a pipeline. Each run tracks:
Checkpoints are pipeline-scoped, adapter-managed state. Extractors and other
components read and update only the offsets or cursors they implement. During a
normal execution, setCheckpoint() updates in-memory state and dirty state is
persisted when execution finalizes; a terminal failure does not automatically
save the last successful record. Later runs load the existing checkpoint unless
an adapter-specific reset option, such as a file extractor’s
resetCheckpoint, is used.
Templates are pre-configured pipeline blueprints for the import and export wizards. They pre-fill wizard steps with common field mappings, source types, and format options.
Built-in import templates:
Built-in export templates:
Custom templates can be registered via plugin options or the TemplateRegistryService.
Scripts are named functions that can modify records at the 18 data-processing stages. Register scripts via plugin options and reference them in pipeline hook definitions. Lifecycle and error stages are observation-only:
| Hook Type | Purpose | Can Modify Records |
|---|---|---|
INTERCEPTOR |
Inline JavaScript for trusted administrators, with restricted globals and a timeout | Yes |
SCRIPT |
Pre-registered TypeScript functions | Yes |
WEBHOOK |
HTTP notifications | No |
EMIT |
Vendure domain events | No |
TRIGGER_PIPELINE |
Start another pipeline | No |
LOG |
Write to pipeline logs | No |
Control how records flow through steps:
{
throughput: {
batchSize: 100, // Process 100 records at a time
concurrency: 4, // Process 4 batches in parallel
rateLimitRps: 10, // Max 10 aggregate load-batch starts per second
pauseOnErrorRate: { threshold: 0.5, intervalSec: 60 }, // React at 50% failures in the rolling 60-second window
drainStrategy: 'BACKOFF', // 'BACKOFF' | 'SHED' | 'QUEUE'
}
}
Batch size is 1-10,000, concurrency is 1-16, and the rate limit is
0-1,000 batch starts per second (0 disables pacing). The error threshold is
greater than 0 and at most 1. Its optional interval is 0.1-3,600 seconds;
BACKOFF defaults to 1 second and QUEUE defaults to 5 seconds.
When loading entities, choose a strategy:
| Strategy | Behavior |
|---|---|
CREATE |
Only create new records; skip if exists |
UPDATE |
Only update existing records; skip if not found |
UPSERT |
Create new or update existing records |
MERGE |
Merge source with existing data |
SOFT_DELETE |
Mark as deleted / logical delete |
HARD_DELETE |
Permanently remove from database |
When loading to Vendure, control channel assignment:
| Strategy | Behavior |
|---|---|
EXPLICIT |
Add to specified channel(s) |
INHERIT |
Inherit channel from parent entity |
MULTI |
Assign to multiple channels |
Records can be validated at multiple points:
Invalid records are quarantined for review and can be:
When a record fails: