@azlib/etl
v0.2.0
Published
A high-performance, framework-agnostic **Extract-Transform-Load (ETL)** pipeline engine for TypeScript and Node.js. Built for resilience, backpressure, atomic batch commits, and zero data loss.
Readme
@azlib/etl
A high-performance, framework-agnostic Extract-Transform-Load (ETL) pipeline engine for TypeScript and Node.js. Built for resilience, backpressure, atomic batch commits, and zero data loss.
Highlights
- Framework Agnostic: Zero external runtime dependencies. Works in Node.js, serverless, edge runtimes, or background workers.
- Zero Data Loss Guarantee:
- Committed Checkpointing: Source cursors and offsets are ONLY committed after the target sink acknowledges successful persistence.
- Dead Letter Queue (DLQ): Failed/corrupted records are captured with complete error diagnostics instead of silently dropped.
- Resumable Pipelines: If interrupted or crashed, resumes execution seamlessly from the last committed checkpoint.
- Replay Utility: Built-in
reprocessDeadLetters()tool to repair and re-inject DLQ records.
- Fluent Pipeline Builder: Declarative
.extract(),.transform(),.pipe(),.load(), and.onError()API. - Backpressure & Streaming: Processes massive datasets in configurable batch sizes with low memory overhead.
- Automatic Retries: Exponential backoff with jitter for transient extract/load errors.
- Graceful Shutdown: Native
AbortSignalsupport.
Installation
pnpm add @azlib/etlQuick Start
import {
createPipeline,
fromArray,
map,
filter,
toArrayLoader,
createInMemoryCheckpointStore,
} from "@azlib/etl";
interface RawUser {
id: number;
name: string;
email: string;
}
interface ProcessedUser {
id: number;
username: string;
email: string;
}
const rawData: RawUser[] = [
{ id: 1, name: "Alice", email: "[email protected]" },
{ id: 2, name: "Bob", email: "[email protected]" },
];
const loader = toArrayLoader<ProcessedUser>();
const pipeline = createPipeline<RawUser, ProcessedUser>({
name: "users-pipeline",
batchSize: 50,
})
.extract(fromArray(rawData))
.pipe(filter((u) => Boolean(u.email)))
.pipe(
map((u) => ({
id: u.id,
username: u.name.toLowerCase(),
email: u.email.toLowerCase(),
})),
)
.load(loader);
const result = await pipeline.run();
console.log(result.metrics);
// { totalExtracted: 2, totalTransformed: 2, totalLoaded: 2, status: 'completed' }Architecture & Zero Data Loss Strategy
[ Data Source ]
│
▼
┌──────────────┐ Extracts in chunks
│ Extractor │ ◄─── Uses checkpoint (cursor/offset)
└──────┬───────┘
│
▼
┌──────────────┐ 1:1, 1:N, or 1:0 (filter)
│ Transformers │ ───► Invalid records ───► [ Dead Letter Queue ]
└──────┬───────┘ │
│ ▼
▼ reprocessDeadLetters()
┌──────────────┐
│ Batch Buffer │ Respects batchSize & concurrency
└──────┬───────┘
│
▼
┌──────────────┐ Retry with exponential backoff
│ Target Sink │ ───► Unrecoverable batch ─► [ Dead Letter Queue ]
└──────┬───────┘
│
▼ (Only AFTER batch load succeeds)
┌──────────────────────┐
│ Commit Checkpoint │ ───► [ Checkpoint Store ]
└──────────────────────┘1. Resumable Pipelines & Checkpoints
When processing millions of records or streaming from databases and message queues, failures can occur midway. Using a CheckpointStore:
import {
createPipeline,
fromCursor,
createFileCheckpointStore,
} from "@azlib/etl";
const checkpointStore = createFileCheckpointStore<string>("./etl-cursor.json");
const pipeline = createPipeline({
name: "sales-sync",
checkpointStore,
})
.extract(
fromCursor(async (cursor, context) => {
// Starts automatically from context.initialCheckpoint!
const page = await fetchDatabasePage({ afterCursor: cursor, limit: 100 });
return {
records: page.items,
nextCursor: page.nextCursor,
hasMore: page.hasNextPage,
};
}),
)
.load(async (batch) => {
await warehouse.insertBatch(batch);
});
await pipeline.run();If the pipeline stops at record 50,000, running it again will resume starting from record 50,001.
2. Dead Letter Queue & Replay
Instead of aborting an entire multi-hour batch job due to a handful of malformed rows, route them to a DeadLetterSink:
import {
createPipeline,
fromArray,
validate,
createFileDeadLetterSink,
reprocessDeadLetters,
} from "@azlib/etl";
const dlq = createFileDeadLetterSink("./failed-records.jsonl");
const pipeline = createPipeline({
name: "import-customers",
deadLetterSink: dlq,
transformErrorStrategy: "dead-letter",
})
.extract(fromArray(customers))
.transform(validate((c) => isValidTaxId(c.taxId), "Invalid Tax ID format"))
.load(async (batch) => {
await database.customers.bulkInsert(batch);
});
await pipeline.run();
// Later: Fix data issue and reprocess failed records!
await reprocessDeadLetters({
sink: dlq,
handler: async (payload, dlqRecord) => {
const sanitized = sanitizeTaxId(payload);
await database.customers.insert(sanitized);
},
});Error Handling Strategies
| Strategy | Behavior |
| :-------------------- | :------------------------------------------------------------------------------------------------------------ |
| 'abort' (default) | Immediately halts execution, rolls back batch, and does NOT advance checkpoint. |
| 'skip' | Discards the offending record/batch, logs a warning, advances checkpoint, and continues. |
| 'dead-letter' | Stores the record and its diagnostic context (error, stack, stage, payload) in the DLQ and continues. |
You can configure error strategies globally or per stage:
createPipeline({
transformErrorStrategy: "dead-letter", // Catch validation errors in DLQ
loadErrorStrategy: "abort", // Never skip downstream database write errors
});API Reference
Extractors
fromArray(items, chunkSize?): Extracts from in-memory arrays or async factories.fromAsyncIterable(iterable, chunkSize?): Extracts from async generators and streams.fromCursor(fetchPage): Checkpoint-aware pagination extractor.
Transformers
map(fn): 1-to-1 sync or async record transformation.filter(predicate): Drop records conditionally.flatMap(fn): Expand 1 record into multiple records.validate(validator, message?): Assert record validity; routes to DLQ on failure.tap(fn): Execute side-effects (e.g. logging/metrics) without altering data.compose(...transformers): Chain multiple transformers into one.
Loaders
toArrayLoader(): In-memory collection loader with.getRecords().batchLoader(fn, options): Batch writer with lifecycle hooks (beforeBatch,afterBatch,flush).
Resilience & DLQ
createInMemoryCheckpointStore(),createFileCheckpointStore(path)createInMemoryDeadLetterSink(),createFileDeadLetterSink(path),createCallbackDeadLetterSink(cb)reprocessDeadLetters({ sink, handler, filter? })executeWithRetry(operation, retryPolicy, signal)
License
MIT © hanhn-dev
