| title | Etl |
|---|---|
| description | Etl protocol schemas |
{/*
ETL (Extract, Transform, Load) Pipeline Protocol - LEVEL 2: Data Engineering
Inspired by modern data integration platforms like Airbyte, Fivetran, and Apache NiFi.
Positioning in 3-Layer Architecture:
-
L1: Simple Sync (automation/sync.zod.ts) - Business users - Sync Salesforce to Sheets
-
L2: ETL Pipeline (THIS FILE) - Data engineers - Aggregate 10 sources to warehouse
-
L3: Enterprise Connector (integration/connector.zod.ts) - System integrators - Full SAP integration
ETL pipelines enable automated data synchronization between systems, transforming
data as it moves from source to destination.
SCOPE: Advanced multi-source, multi-stage transformations.
Supports complex operations: joins, aggregations, filtering, custom SQL.
Use ETL Pipeline when:
-
Combining data from multiple sources
-
Need aggregations, joins, transformations
-
Building data warehouses or analytics platforms
-
Complex data transformations required
Examples:
-
Sales data from Salesforce + Marketing from HubSpot → Data Warehouse
-
Multi-region databases → Consolidated reporting
-
Legacy system migration with transformation
When to downgrade:
- Simple 1:1 sync → Use Simple Sync
When to upgrade:
- Need full connector lifecycle (auth, webhooks, rate limits) → Use Enterprise Connector
See also: ./sync.zod.ts for Level 1 (simple sync)
See also: [../integration/connector.zod.ts](/docs/references/integration/connector) for Level 3 (enterprise integration)
- Data Warehouse Population
-
Extract from multiple operational systems
-
Transform to analytical schema
-
Load into data warehouse
- System Integration
-
Sync data between CRM and Marketing Automation
-
Keep product catalog synchronized across e-commerce platforms
-
Replicate data for backup/disaster recovery
- Data Migration
-
Move data from legacy systems to modern platforms
-
Consolidate data from multiple sources
-
Split monolithic databases into microservices
See also: https://airbyte.com/
See also: https://docs.fivetran.com/
See also: https://nifi.apache.org/
@example
const salesforceToDB: ETLPipeline = \{
name: 'salesforce_to_postgres',
label: 'Salesforce Accounts to PostgreSQL',
source: \{
type: 'api',
connector: 'salesforce',
config: \{ object: 'Account' \}
\},
destination: \{
type: 'database',
connector: 'postgres',
config: \{ table: 'accounts' \}
\},
transformations: [
\{ type: 'map', config: \{ 'Name': 'account_name' \} \}
],
schedule: '0 2 * * *' // Daily at 2 AM
\}import { ETLDestination, ETLEndpointType, ETLPipeline, ETLPipelineRun, ETLRunStatus, ETLSource, ETLSyncMode, ETLTransformation, ETLTransformationType } from '@objectstack/spec/automation';
import type { ETLDestination, ETLEndpointType, ETLPipeline, ETLPipelineRun, ETLRunStatus, ETLSource, ETLSyncMode, ETLTransformation, ETLTransformationType } from '@objectstack/spec/automation';
// Validate data
const result = ETLDestination.parse(data);| Property | Type | Required | Description |
|---|---|---|---|
| type | Enum<'database' | 'api' | 'file' | 'stream' | 'object' | 'warehouse' | 'storage' | 'spreadsheet'> |
✅ | Destination type |
| connector | string |
optional | Connector ID |
| config | Record<string, any> |
✅ | Destination configuration |
| writeMode | Enum<'append' | 'overwrite' | 'upsert' | 'merge'> |
✅ | How to write data |
| primaryKey | string[] |
optional | Primary key fields |
databaseapifilestreamobjectwarehousestoragespreadsheet
| Property | Type | Required | Description |
|---|---|---|---|
| name | string |
✅ | Pipeline identifier (snake_case) |
| label | string |
optional | Pipeline display name |
| description | string |
optional | Pipeline description |
| source | { type: Enum<'database' | 'api' | 'file' | 'stream' | 'object' | 'warehouse' | 'storage' | 'spreadsheet'>; connector?: string; config: Record<string, any>; incremental?: object } |
✅ | Data source |
| destination | { type: Enum<'database' | 'api' | 'file' | 'stream' | 'object' | 'warehouse' | 'storage' | 'spreadsheet'>; connector?: string; config: Record<string, any>; writeMode?: Enum<'append' | 'overwrite' | 'upsert' | 'merge'>; … } |
✅ | Data destination |
| transformations | { name?: string; type: Enum<'map' | 'filter' | 'aggregate' | 'join' | 'script' | 'lookup' | 'split' | 'merge' | 'normalize' | 'deduplicate'>; config: Record<string, any>; continueOnError?: boolean }[] |
optional | Transformation pipeline |
| syncMode | Enum<'full' | 'incremental' | 'cdc'> |
optional | Sync mode |
| schedule | string | { dialect: Enum<'cel' | 'js' | 'cron' | 'template'>; source?: string; ast?: any; meta?: object } |
optional | Cron schedule expression |
| enabled | boolean |
optional | Pipeline enabled status |
| retry | { maxAttempts?: integer; backoffMs?: integer } |
optional | Retry configuration |
| notifications | { onSuccess?: string[]; onFailure?: string[] } |
optional | Notification settings |
| tags | string[] |
optional | Pipeline tags |
| metadata | Record<string, any> |
optional | Custom metadata |
| Property | Type | Required | Description |
|---|---|---|---|
| id | string |
✅ | Run identifier |
| pipelineName | string |
✅ | Pipeline name |
| status | Enum<'pending' | 'running' | 'succeeded' | 'failed' | 'cancelled' | 'timeout'> |
✅ | Run status |
| startedAt | string |
✅ | Start time |
| completedAt | string |
optional | Completion time |
| durationMs | number |
optional | Duration in ms |
| stats | { recordsRead: integer; recordsWritten: integer; recordsErrored: integer; bytesProcessed: integer } |
optional | Run statistics |
| error | { message: string; code?: string; details?: any } |
optional | Error information |
| logs | string[] |
optional | Execution logs |
pendingrunningsucceededfailedcancelledtimeout
| Property | Type | Required | Description |
|---|---|---|---|
| type | Enum<'database' | 'api' | 'file' | 'stream' | 'object' | 'warehouse' | 'storage' | 'spreadsheet'> |
✅ | Source type |
| connector | string |
optional | Connector ID |
| config | Record<string, any> |
✅ | Source configuration |
| incremental | { enabled: boolean; cursorField: string; cursorValue?: any } |
optional | Incremental extraction config |
fullincrementalcdc
| Property | Type | Required | Description |
|---|---|---|---|
| name | string |
optional | Transformation name |
| type | Enum<'map' | 'filter' | 'aggregate' | 'join' | 'script' | 'lookup' | 'split' | 'merge' | 'normalize' | 'deduplicate'> |
✅ | Transformation type |
| config | Record<string, any> |
✅ | Transformation config |
| continueOnError | boolean |
✅ | Continue on error |
mapfilteraggregatejoinscriptlookupsplitmergenormalizededuplicate