Build data pipelines using the visual drag-and-drop editor.
Pipeline Management - View and manage all your data pipelines
The editor has three main areas:
Simple Mode - Step-by-step list view for building pipelines
Workflow Mode - Visual drag-and-drop canvas with node palette
Every pipeline needs a trigger to define how it starts.
Extract steps pull data from sources.
httpApi) - REST API endpoints with pagination, authentication, and retry supportcsv, json, xml, xlsx) - Parse Data Hub uploads or the inline fields supported by the selected formatgraphql) - External GraphQL endpointsvendureQuery) - Vendure entity dataTransform steps modify records.
Load steps create or update Vendure entities or send data externally.
productUpsert) - Create/update productsvariantUpsert) - Create/update product variantscustomerUpsert) - Create/update customerscollectionUpsert) - Create/update collectionspromotionUpsert) - Create/update promotionsstockAdjust) - Adjust inventory levelsorderNote) - Add notes to ordersorderTransition) - Change order statesrestPost) - Send data to external APIsConnections show the data flow direction. A record must pass through connected steps in order.
Click a step to open its configuration panel:
Each adapter has specific settings. See Reference for details.
Extract and Validate steps can bind to a named version from Data Hub > Schemas. A strict or backward-compatible mismatch rejects the record; permissive mode records a warning and accepts it. Step tests and dry runs use the same binding as live execution. See Schema Registry.
Route steps split data flow based on conditions:
Example: Route products by category:
Before running, validate your pipeline:
Fix all issues before saving.
Removing its deployed definition releases the persisted pipeline to Dashboard ownership on the next API startup without deleting revisions or run history. Review its schedules, triggers, references, and current published state before editing, enabling, or deleting the released pipeline.
Some pipelines accept input parameters:
| State | Description |
|---|---|
DRAFT |
Editable working definition; runs the previous published revision when one exists and the pipeline is enabled |
REVIEW |
Submitted working definition; runs the previous published revision when one exists and the pipeline is enabled |
PUBLISHED |
Working definition matches the selected published revision |
ARCHIVED |
Retired definition; cannot run |
Lifecycle status and the enabled switch are separate. A pipeline must be both
enabled and have a selected published revision before it can run, and it must
not be archived. Draft and review edits never enter production execution until
they are published. Pipeline codes become immutable after the first
publication so webhook URLs and cross-pipeline dependencies remain stable.
Execution state such as RUNNING, COMPLETED, or FAILED belongs to an
individual pipeline run, not the pipeline.
Note: Deleting a pipeline removes all run history.
Hooks allow you to execute custom code at specific stages of pipeline execution. Use hooks to modify data, send notifications, trigger other pipelines, or integrate with external systems.
Data Processing Stages:
| Stage | When It Runs | Can Modify Records |
|---|---|---|
BEFORE_EXTRACT |
Before data extraction | Yes (seed records) |
AFTER_EXTRACT |
After data is extracted | Yes |
BEFORE_TRANSFORM |
Before transformation | Yes |
AFTER_TRANSFORM |
After transformation | Yes |
BEFORE_VALIDATE |
Before validation | Yes |
AFTER_VALIDATE |
After validation | Yes |
BEFORE_ENRICH |
Before enrichment | Yes |
AFTER_ENRICH |
After enrichment | Yes |
BEFORE_ROUTE |
Before routing | Yes |
AFTER_ROUTE |
After routing | Yes |
BEFORE_LOAD |
Before loading to Vendure | Yes |
AFTER_LOAD |
After loading | Yes |
BEFORE_EXPORT |
Before file export | Yes |
AFTER_EXPORT |
After file export | Yes |
BEFORE_FEED |
Before feed generation | Yes |
AFTER_FEED |
After feed generation | Yes |
BEFORE_SINK |
Before search indexing | Yes |
AFTER_SINK |
After search indexing | Yes |
Lifecycle Stages (observe-only — WEBHOOK, EMIT, LOG, TRIGGER_PIPELINE only, no INTERCEPTOR/SCRIPT):
| Stage | When It Runs | Supported Hook Types |
|---|---|---|
PIPELINE_STARTED |
Pipeline execution begins | WEBHOOK, EMIT, LOG, TRIGGER_PIPELINE |
PIPELINE_COMPLETED |
Pipeline finishes successfully | WEBHOOK, EMIT, LOG, TRIGGER_PIPELINE |
PIPELINE_FAILED |
Pipeline fails | WEBHOOK, EMIT, LOG, TRIGGER_PIPELINE |
ON_ERROR |
When an error occurs | WEBHOOK, EMIT, LOG, TRIGGER_PIPELINE |
ON_RETRY |
When a record is retried | WEBHOOK, EMIT, LOG, TRIGGER_PIPELINE |
ON_DEAD_LETTER |
When a record is sent to dead letter queue | WEBHOOK, EMIT, LOG, TRIGGER_PIPELINE |
Interceptors run JavaScript code that can modify the records array:
.hooks({
AFTER_EXTRACT: [{
type: 'INTERCEPTOR',
name: 'Add metadata',
code: `
return records.map(r => ({
...r,
extractedAtEpochMs: Date.now(),
source: 'supplier-api',
}));
`,
failOnError: false, // Don't fail pipeline if hook fails
timeout: 5000, // 5 second timeout
}],
BEFORE_LOAD: [{
type: 'INTERCEPTOR',
name: 'Filter invalid',
code: `
return records.filter(r => r.sku && r.name);
`,
}],
})
Modify records before search indexing (Meilisearch, Elasticsearch, etc.):
.hooks({
BEFORE_SINK: [{
type: 'INTERCEPTOR',
name: 'Enrich for search',
code: `
return records.map(r => ({
...r,
searchText: [r.name, r.sku, r.description].filter(Boolean).join(' ').toLowerCase(),
facetTags: (r.tags || '').split(',').map(t => t.trim()).filter(Boolean),
boostScore: r.featured ? 1.5 : 1.0,
}));
`,
}],
})
Transform records before CSV/JSON export:
.hooks({
BEFORE_EXPORT: [{
type: 'INTERCEPTOR',
name: 'Format for export',
code: `
return records.map(r => ({
...r,
price: (r.price / 100).toFixed(2),
createdAtEpochMs: Date.parse(r.createdAt),
}));
`,
}],
})
Available in interceptor code:
records - A deep-cloned current record arraycontext - A deep-cloned hook context with pipelineId, runId, stage, and recordsArray, Object, String, Number, JSON, and Math membersDate.now() and Date.parse(); Date is not available as a constructorisNaN, isFinite, and sandboxed console.log/warn/errorNetwork access, module loading, new Date(), promises, timers, and async syntax
are not supported. Use a registered SCRIPT hook for trusted TypeScript logic
that needs host APIs.
Script hooks reference pre-registered TypeScript functions. Register scripts via plugin options (recommended) or imperatively:
Via Plugin Options (Recommended):
DataHubPlugin.init({
scripts: {
'addSegment': async (records, context, args) => {
const threshold = args?.threshold || 1000;
return records.map(r => ({
...r,
segment: r.totalSpent > threshold ? 'vip' : 'standard',
}));
},
'validateRequired': async (records, context) => {
return records.filter(r => r.sku && r.name && r.price > 0);
},
'enrichWithTimestamp': async (records, context) => {
return records.map(r => ({
...r,
importedAt: Date.now(),
pipelineRun: context.runId,
}));
},
},
})
Via Service Injection:
import { HookService } from '@oronts/vendure-data-hub-plugin';
@VendurePlugin({ imports: [DataHubPlugin] })
export class MyPlugin implements OnModuleInit {
constructor(private hookService: HookService) {}
onModuleInit() {
this.hookService.registerScript('addSegment', async (records, context, args) => {
const threshold = args?.threshold || 1000;
return records.map(r => ({
...r,
segment: r.totalSpent > threshold ? 'vip' : 'standard',
}));
});
}
}
Tip: When registering scripts via
HookService.registerScript()in a NestJS service, your scripts can access any injected service (database, external APIs, Vendure services) through JavaScript closures. See the Developer Guide for examples.
Use in pipeline:
.hooks({
AFTER_TRANSFORM: [{
type: 'SCRIPT',
scriptName: 'addSegment',
args: { threshold: 5000 },
}],
AFTER_EXTRACT: [{
type: 'SCRIPT',
scriptName: 'validateRequired',
}],
})
Register a script for search index enrichment:
DataHubPlugin.init({
scripts: {
'buildSearchAttributes': async (records, context) => {
return records.map(r => ({
...r,
searchText: [r.name, r.sku, r.description]
.filter(Boolean).join(' ').toLowerCase(),
facetCategories: r.categories?.map(c => c.name) || [],
}));
},
},
})
// Use in pipeline:
.hooks({
BEFORE_SINK: [{ type: 'SCRIPT', scriptName: 'buildSearchAttributes' }],
})
Send HTTP notifications to external systems:
.hooks({
PIPELINE_COMPLETED: [{
type: 'WEBHOOK',
url: 'https://slack.example.com/webhook',
headers: { 'Content-Type': 'application/json' },
}],
PIPELINE_FAILED: [{
type: 'WEBHOOK',
url: 'https://pagerduty.example.com/alert',
secretCode: 'webhook-signing-secret', // HMAC signing Secret Code
signatureHeader: 'X-Signature',
retryConfig: {
maxAttempts: 5,
initialDelayMs: 1000,
backoffMultiplier: 2,
},
}],
})
Emit Vendure domain events:
.hooks({
PIPELINE_COMPLETED: [{
type: 'EMIT',
event: 'ProductSyncCompleted',
}],
})
Start another pipeline with the current records:
.hooks({
AFTER_LOAD: [{
type: 'TRIGGER_PIPELINE',
pipelineCode: 'reindex-search',
triggerKey: 'hook',
}],
})
The action creates and queues a pending child run asynchronously. The parent
does not wait for child completion or inherit the child outcome. A
failOnError setting covers only immediate child creation and queue-request
failure.
Test configured observation actions without running the full pipeline:
The Hooks page is an observability and side-effect test surface. WEBHOOK,
EMIT, TRIGGER_PIPELINE, and LOG actions execute. INTERCEPTOR and
SCRIPT actions require the real record-processing lifecycle and are reported
as skipped; use a pipeline dry run to inspect their record modifications.
failOnError: false for non-critical hooksThe Data Hub provides guided wizards for creating import and export pipelines:
Wizards offer pre-configured templates for common scenarios:
Import Templates: REST API Sync, JSON Import, Magento CSV, XML Feed, ERP Inventory, CRM Customer Sync Export Templates: Product XML/CSV/JSON, Order Analytics/CSV, Customer GDPR/CSV
Marketplace feeds are configured from Data Hub > Feeds, where Google, Meta, and Amazon formats use their dedicated generators and catalog data contract.
Custom templates registered via plugin options or connectors appear alongside built-in templates.
Open Version history from a pipeline detail page to inspect revisions, run counts, and the most recent run outcome for each published revision. The dialog shows the latest 20 revisions. The Admin API accepts an integer limit from 1 to 500 and defaults to 50.
Test pipelines without persisting changes:
Dry run executes extract, transform, validate, route, and loader simulation paths without loader writes. ENRICH, EXPORT, FEED, SINK, and GATE side effects are deliberately not executed; the result marks each such step as skipped and preserves its input records. Use a controlled staging run to verify external delivery, credentials, approval behavior, and production write constraints.
product-import-daily not pipeline-1inventory-sync-hourlyerp-product-synccontinue, stop, or dead-letter