DtPipe.Core
1.2.3
See the version list below for details.
dotnet add package DtPipe.Core --version 1.2.3
NuGet\Install-Package DtPipe.Core -Version 1.2.3
<PackageReference Include="DtPipe.Core" Version="1.2.3" />
<PackageVersion Include="DtPipe.Core" Version="1.2.3" />
<PackageReference Include="DtPipe.Core" />
paket add DtPipe.Core --version 1.2.3
#r "nuget: DtPipe.Core, 1.2.3"
#:package DtPipe.Core@1.2.3
#addin nuget:?package=DtPipe.Core&version=1.2.3
#tool nuget:?package=DtPipe.Core&version=1.2.3
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
└── PipelineEngine.cs # Headless pipeline engine for library use
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
var engine = new PipelineEngine(logger);
long rowsWritten = await engine.RunAsync(
reader,
writer,
pipeline: transformers, // optional
batchSize: 50_000,
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 (>= 22.1.0)
- Apache.Arrow.Serialization (>= 1.2.3)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.5)
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 |