DtPipe.Core
1.8.2
See the version list below for details.
dotnet add package DtPipe.Core --version 1.8.2
NuGet\Install-Package DtPipe.Core -Version 1.8.2
<PackageReference Include="DtPipe.Core" Version="1.8.2" />
<PackageVersion Include="DtPipe.Core" Version="1.8.2" />
<PackageReference Include="DtPipe.Core" />
paket add DtPipe.Core --version 1.8.2
#r "nuget: DtPipe.Core, 1.8.2"
#:package DtPipe.Core@1.8.2
#addin nuget:?package=DtPipe.Core&version=1.8.2
#tool nuget:?package=DtPipe.Core&version=1.8.2
DtPipe.Core
The core pipeline engine for DtPipe. Contains all abstractions, models, and pipeline logic with no external dependencies beyond Microsoft.Extensions.Logging.Abstractions.
Suitable for use as a standalone NuGet package in custom ETL pipelines.
Package
<PackageReference Include="DtPipe.Core" Version="1.0.0" />
What's Inside
DtPipe.Core/
├── Abstractions/ # Core interfaces (row and columnar)
├── Attributes/ # [ComponentOption] attribute
├── Dialects/ # ISqlDialect, SQL generation helpers
├── Helpers/ # Shared utility helpers
├── Infrastructure/ # Arrow type mapping, schema factory, row↔columnar bridges
│ └── Arrow/
├── Models/ # Shared data models (PipeColumnInfo, etc.)
├── Options/ # IOptionSet, IQueryAwareOptions, IKeyAwareOptions, OptionsRegistry
├── Pipelines/ # PipelineExecutionPlan, PipelineSegment, DAG orchestration
│ └── Dag/
├── Security/ # SQL query validator
└── Validation/ # Schema and constraint validators
Key Abstractions
| Interface | Purpose |
|---|---|
IStreamReader |
Reads rows as async batches from a source |
IStreamReaderFactory |
Creates IStreamReader instances |
IColumnarStreamReader |
Reads Apache Arrow RecordBatch streams (zero-copy columnar path) |
IDataWriter |
Base contract for all writers |
IRowDataWriter |
Writes object?[] rows to a destination |
IColumnarDataWriter |
Writes Arrow RecordBatch directly (no row conversion) |
IDataWriterFactory |
Creates IDataWriter instances |
IDataTransformer |
Transforms a batch of rows in the pipeline |
IDataTransformerFactory |
Creates IDataTransformer instances |
IColumnarTransformer |
Columnar (Arrow-level) variant of IDataTransformer |
IMultiRowTransformer |
Transformer that may produce multiple rows per input row |
IColumnTypeInferenceCapable |
Reader that supports --auto-column-types inference |
IProviderDescriptor<T> |
Describes a provider: name, options type, factory method |
ISchemaInspector |
Introspects target schema |
ISchemaMigrator |
Applies schema migrations (auto-migrate) |
ISqlDialect |
Generates provider-specific DDL/DML SQL |
Key Options Interfaces
| Interface | Purpose |
|---|---|
IOptionSet |
Base contract for all Options classes |
IQueryAwareOptions |
Implement on reader options to receive the global --query |
IKeyAwareOptions |
Implement on writer options to receive the global --key |
Implementing a Custom Reader
// 1. Define options
public record MyReaderOptions : IProviderOptions, IQueryAwareOptions
{
public static string Prefix => "mydb";
public static string DisplayName => "MyDB Reader";
public string? Query { get; set; }
}
// 2. Implement IStreamReader
public class MyStreamReader : IStreamReader
{
public IReadOnlyList<PipeColumnInfo>? Columns { get; private set; }
public Task OpenAsync(CancellationToken ct) { ... }
public IAsyncEnumerable<ReadOnlyMemory<object?[]>> ReadBatchesAsync(int batchSize, CancellationToken ct) { ... }
public ValueTask DisposeAsync() { ... }
}
// 3. Implement IProviderDescriptor<IStreamReader>
public class MyReaderDescriptor : IProviderDescriptor<IStreamReader>
{
public string ProviderName => "mydb";
public Type OptionsType => typeof(MyReaderOptions);
public bool RequiresQuery => true;
public bool CanHandle(string cs) => cs.StartsWith("mydb:");
public IStreamReader Create(string cs, object options, IServiceProvider sp)
=> new MyStreamReader(cs, ((MyReaderOptions)options).Query!);
}
Running the Pipeline Directly
Headless execution goes through DagOrchestrator — even a single-branch (linear)
run goes through it once, giving uniform cancellation and channel wiring:
var orchestrator = new DagOrchestrator(logger, channelRegistry, readerFactories);
var job = new JobDefinition { Input = "csv:in.csv", Output = "parquet:out.parquet" };
var dag = new JobDagDefinition
{
// BranchDefinition.FromJob is the one projection from a job to a branch; building the
// record by hand is how four call sites drifted apart.
Branches = new[] { BranchDefinition.FromJob("main", job) }
};
int exitCode = await orchestrator.ExecuteAsync(dag, (branch, ctx, ct) =>
{
// Open the branch's reader, run transformers, write to its writer.
return RunBranchAsync(branch, ctx, ct);
}, cancellationToken);
License
MIT
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net10.0 is compatible. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net10.0
- Apache.Arrow (>= 23.0.0)
- DtPipe.Arrow.Serialization (>= 1.8.2)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.11)
NuGet packages (4)
Showing the top 4 NuGet packages that depend on DtPipe.Core:
| Package | Downloads |
|---|---|
|
DtPipe.Adapters
Database and file adapters for DtPipe.Core (PostgreSQL, Oracle, SQL Server, SQLite, DuckDB, CSV, Parquet, JsonL, Apache Arrow). |
|
|
DtPipe.Processors
Stream processors for DtPipe pipelines (DuckDB SQL, Merge/UNION ALL). |
|
|
DtPipe.Transformers
Row and columnar data transformers for DtPipe pipelines (expand, filter, compute, anonymize, fake data generation). |
|
|
DtPipe.Adapters.Shared
Package Description |
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 1.9.0 | 0 | 10/2/2026 |
| 1.8.2 | 127 | 9/14/2026 |
| 1.8.1 | 129 | 9/10/2026 |
| 1.8.0 | 123 | 9/9/2026 |
| 1.7.0 | 135 | 9/5/2026 |
| 1.6.0 | 136 | 8/27/2026 |
| 1.5.0 | 134 | 8/6/2026 |
| 1.4.3 | 162 | 7/7/2026 |
| 1.4.2 | 160 | 6/20/2026 |
| 1.4.1 | 157 | 6/19/2026 |
| 1.4.0 | 167 | 6/17/2026 |
| 1.3.4 | 155 | 6/15/2026 |
| 1.3.3 | 162 | 6/14/2026 |
| 1.3.2 | 160 | 6/9/2026 |
| 1.3.1 | 163 | 6/8/2026 |
| 1.3.0 | 187 | 5/8/2026 |
| 1.2.6 | 182 | 4/29/2026 |
| 1.2.5 | 188 | 4/25/2026 |
| 1.2.4 | 194 | 4/14/2026 |
| 1.2.3 | 199 | 4/14/2026 |