RtFlow.Pipelines.Hosting
1.0.0
dotnet add package RtFlow.Pipelines.Hosting --version 1.0.0
NuGet\Install-Package RtFlow.Pipelines.Hosting -Version 1.0.0
<PackageReference Include="RtFlow.Pipelines.Hosting" Version="1.0.0" />
<PackageVersion Include="RtFlow.Pipelines.Hosting" Version="1.0.0" />
<PackageReference Include="RtFlow.Pipelines.Hosting" />
paket add RtFlow.Pipelines.Hosting --version 1.0.0
#r "nuget: RtFlow.Pipelines.Hosting, 1.0.0"
#:package RtFlow.Pipelines.Hosting@1.0.0
#addin nuget:?package=RtFlow.Pipelines.Hosting&version=1.0.0
#tool nuget:?package=RtFlow.Pipelines.Hosting&version=1.0.0
RtFlow.Pipelines.Hosting
RtFlow.Pipelines.Hosting provides seamless integration with ASP.NET Core and .NET Generic Host, enabling lifecycle management of data processing pipelines as hosted services.
Installation
# Install core package first
dotnet add package RtFlow.Pipelines.Core
# Then add hosting integration
dotnet add package RtFlow.Pipelines.Hosting
Key Features
- 🏠 ASP.NET Core Integration - Register pipelines as hosted services
- 🔄 Lifecycle Management - Automatic startup, shutdown, and graceful termination
- ⚙️ Configuration Support - Configure pipelines via appsettings.json
- 📊 Health Checks - Built-in health check integration
- 🔍 Observability - Logging and metrics integration with .NET diagnostics
- 🎛️ Background Processing - Long-running pipeline services
Quick Start
Basic Hosted Pipeline
using RtFlow.Pipelines.Hosting;
var builder = WebApplication.CreateBuilder(args);
// Register a pipeline as a hosted service
builder.Services.AddPipelineHostedService<DataProcessor>(serviceProvider =>
{
return PipelineBuilder
.Create<DataItem>()
.Transform(item => ProcessItem(item))
.ForEach(result => SaveResult(result))
.Build();
});
var app = builder.Build();
app.Run();
Pipeline Factory Service
public class DataProcessingService : PipelineHostedService<DataItem>
{
private readonly ILogger<DataProcessingService> _logger;
private readonly IDataSource _dataSource;
public DataProcessingService(
ILogger<DataProcessingService> logger,
IDataSource dataSource)
{
_logger = logger;
_dataSource = dataSource;
}
protected override IPipeline<DataItem> CreatePipeline()
{
return PipelineBuilder
.Create<DataItem>()
.Transform(async item => await EnrichData(item))
.Filter(item => item.IsValid)
.Batch(100, TimeSpan.FromSeconds(30))
.Transform(batch => ProcessBatch(batch))
.SideEffect(result => _logger.LogInformation("Processed batch: {Count}", result.Count))
.Build();
}
protected override async IAsyncEnumerable<DataItem> GetDataSource(
[EnumeratorCancellation] CancellationToken cancellationToken)
{
await foreach (var item in _dataSource.StreamAsync(cancellationToken))
{
yield return item;
}
}
}
// Register in Program.cs
builder.Services.AddHostedService<DataProcessingService>();
Configuration
appsettings.json Configuration
{
"PipelineOptions": {
"MaxConcurrency": 10,
"BoundedCapacity": 1000,
"EnableMetrics": true,
"GracefulShutdownTimeout": "00:00:30"
}
}
Dependency Injection
// Configure pipeline options
builder.Services.Configure<PipelineOptions>(
builder.Configuration.GetSection("PipelineOptions"));
// Register dependencies
builder.Services.AddScoped<IDataProcessor, DataProcessor>();
builder.Services.AddScoped<IDataRepository, DataRepository>();
// Register pipeline with dependencies
builder.Services.AddPipelineHostedService<OrderData>(serviceProvider =>
{
var processor = serviceProvider.GetRequiredService<IDataProcessor>();
var repository = serviceProvider.GetRequiredService<IDataRepository>();
return PipelineBuilder
.Create<OrderData>()
.Transform(order => processor.ProcessOrder(order))
.ForEach(result => repository.SaveResult(result))
.Build();
});
Health Checks
builder.Services.AddHealthChecks()
.AddPipelineHealthCheck<DataProcessingService>("data-pipeline");
// Health check endpoint
app.MapHealthChecks("/health");
Graceful Shutdown
The hosting integration automatically handles graceful shutdown:
public class DataProcessingService : PipelineHostedService<DataItem>
{
protected override async Task OnStoppingAsync(CancellationToken cancellationToken)
{
_logger.LogInformation("Pipeline is shutting down gracefully...");
// Custom cleanup logic
await FlushPendingData();
await base.OnStoppingAsync(cancellationToken);
}
}
Background Processing Patterns
Continuous Processing
public class ContinuousProcessorService : PipelineHostedService<QueueMessage>
{
protected override async IAsyncEnumerable<QueueMessage> GetDataSource(
[EnumeratorCancellation] CancellationToken cancellationToken)
{
while (!cancellationToken.IsCancellationRequested)
{
var messages = await _messageQueue.ReceiveAsync(cancellationToken);
foreach (var message in messages)
{
yield return message;
}
}
}
}
Scheduled Processing
public class ScheduledProcessorService : BackgroundService
{
private readonly IPipeline<DataItem> _pipeline;
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
using var timer = new PeriodicTimer(TimeSpan.FromHours(1));
while (!stoppingToken.IsCancellationRequested &&
await timer.WaitForNextTickAsync(stoppingToken))
{
var data = await LoadBatchData();
await _pipeline.ExecuteAsync(data, stoppingToken);
}
}
}
Requirements
- .NET 8.0 or later
- RtFlow.Pipelines.Core package
- Microsoft.Extensions.Hosting package
Documentation
For complete documentation and examples, visit the main project repository.
License
This project is licensed under the MIT License - see the LICENSE file for details.
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net8.0 is compatible. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 was computed. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 was computed. 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. |
-
net8.0
- Microsoft.Extensions.Hosting.Abstractions (>= 9.0.5)
- RtFlow.Pipelines.Core (>= 1.0.0)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|
Hosting package for ASP.NET Core integration with dependency injection and lifecycle management.