This is a viewer only at the moment see the article on how this works.
To update the preview hit Ctrl-Alt-R (or ⌘-Alt-R on Mac) or Enter to refresh. The Save icon lets you save the markdown file to disk
This is a preview from the server running through my markdig pipeline
Friday, 12 December 2025
Sisään **Osa 1: Tuli ja älä Aikamoista Unohda**Tutkimme teoriaa, jonka mukaan hetkellinen teloitus - rajatut, yksityiset, torjuttavat asynkkiset työvirrat, jotka muistavat juuri sen verran, että niistä on hyötyä, ja sitten ne haihtuvat.
Tämä artikkeli muuttaa mallin uudelleenkäytettäväksi kirjastoksi, jonka voit pudottaa mihin tahansa .net-projektiin.
Kirjasto on jaettu hyvin tekaistuihin tiedostoihin:
Tiedoston tarkoitus
|------|---------|
| EphemeralOptions.cs Konfiguraatio (valuutta, ikkunan koko, elinikä, signaalit)
| EphemeralOperation.cs Sisäisen toiminnan seuranta signaalituella
| Snapshots.cs Kuluttajalle altistuvat muuttumattomat kuva-aineistot
| Signaalit.cs Signaalitapahtumat, levinneisyys, rajoitteet ja maailmanlaajuinen SignalSink
| EphemeralIdGenerator.cs Nopea XxHash64-pohjainen ID-sukupolvi
| ConcurrencyGates.cs Kiinteä ja säädettävissä oleva valuutta rajoittamassa
| StringPatternMatcher.cs Glob-tyylinen malli, joka vastaa signaalin suodatusta
| Rinnakkaiset ephemeral.cs Staattinen laajennusmenetelmä (EphemeralForEachAsync) |
| EphemeralWorkCoordinator.cs Pitkäikäinen työjonokoordinaattori
| EphemeralKeyedWork Coordinator.cs Per-avain peräkkäinen toteutus reilulla aikataululla
| EphemeralResultCoordinator.cs Tulosten sitomista koordinoivan koordinaattorin variantti
| SignalDispatcher.cs Async-signaalin reititys kuvion kanssa
| RiippuvuusInjektointi.cs DI-laajennusmenetelmät ja tehtaan toteutustavat
| Esimerkkejä/SignalingHttpClient.cs Näytteen hienorakeinen signaaliemissio HTTP-puheluissa
Ja kattavat testit kattaa kaikki reunatapaukset.
Tämän korvaamme:
// ❌ Before: Fire-and-forget black hole
_ = Task.Run(() => ProcessAsync(item));
// No visibility. No debugging. No idea if it worked.
// ❌ Or: Blocking everything
await ProcessAsync(item); // Hope you like waiting...
Ja mitä me rakennamme:
// ✅ After: Trackable, bounded, debuggable
await coordinator.EnqueueAsync(item);
// Instant visibility
Console.WriteLine($"Pending: {coordinator.PendingCount}");
Console.WriteLine($"Active: {coordinator.ActiveCount}");
Console.WriteLine($"Failed: {coordinator.TotalFailed}");
// Full operation history
var snapshot = coordinator.GetSnapshot();
var failures = coordinator.GetFailed();
Sama async-suoritus. Täydellinen havainnoitavuus. Käyttäjätietoja ei säilytetä.
Yleisin kaava - rekisteröi koordinaattori DI:hen ja ruiskuta se:
// Program.cs
services.AddEphemeralWorkCoordinator<TranslationRequest>(
async (request, ct) => await TranslateAsync(request, ct),
new EphemeralOptions { MaxConcurrency = 8 });
// Your service
public class TranslationService(EphemeralWorkCoordinator<TranslationRequest> coordinator)
{
public async Task TranslateAsync(TranslationRequest request)
{
await coordinator.EnqueueAsync(request);
// Returns immediately - work happens in background
}
public object GetStatus() => new
{
pending = coordinator.PendingCount,
active = coordinator.ActiveCount,
completed = coordinator.TotalCompleted,
failed = coordinator.TotalFailed
};
}
┌─────────────────────────────────────────────────────────────────┐
│ DECISION TREE │
├─────────────────────────────────────────────────────────────────┤
│ │
│ Processing a collection once? │
│ └─► EphemeralForEachAsync<T> (ParallelEphemeral.cs) │
│ │
│ Need a long-lived queue that accepts items over time? │
│ └─► EphemeralWorkCoordinator<T> │
│ │
│ Need per-entity ordering (user commands, tenant jobs)? │
│ └─► EphemeralKeyedWorkCoordinator<TKey, T> │
│ │
│ Need to capture results (fingerprints, summaries)? │
│ └─► EphemeralResultCoordinator<TInput, TResult> │
│ │
│ Need multiple coordinators with different configs? │
│ └─► IEphemeralCoordinatorFactory<T> (like IHttpClientFactory) │
│ │
│ Need dynamic concurrency adjustment at runtime? │
│ └─► Set EnableDynamicConcurrency = true, call SetMaxConcurrency│
│ │
└─────────────────────────────────────────────────────────────────┘
From EphemeralOptions.cs:
public sealed class EphemeralOptions
{
// Concurrency control
public int MaxConcurrency { get; init; } = Environment.ProcessorCount;
public int MaxConcurrencyPerKey { get; init; } = 1;
public bool EnableDynamicConcurrency { get; init; } = false;
// Window management
public int MaxTrackedOperations { get; init; } = 200;
public TimeSpan? MaxOperationLifetime { get; init; } = TimeSpan.FromMinutes(5);
// Fair scheduling (keyed coordinator)
public bool EnableFairScheduling { get; init; } = false;
public int FairSchedulingThreshold { get; init; } = 10;
// Signal-reactive processing
public IReadOnlySet<string>? CancelOnSignals { get; init; }
public IReadOnlySet<string>? DeferOnSignals { get; init; }
public int MaxDeferAttempts { get; init; } = 10;
public TimeSpan DeferCheckInterval { get; init; } = TimeSpan.FromMilliseconds(100);
// Signal infrastructure
public SignalSink? Signals { get; init; }
public SignalConstraints? SignalConstraints { get; init; }
public Action<SignalEvent>? OnSignal { get; init; }
// Async signal handling
public Func<SignalEvent, CancellationToken, Task>? OnSignalAsync { get; init; }
public int MaxConcurrentSignalHandlers { get; init; } = 4;
public int MaxQueuedSignals { get; init; } = 1000;
// Observability
public Action<IReadOnlyCollection<EphemeralOperationSnapshot>>? OnSample { get; init; }
}
SetMaxConcurrency() - käyttää kustomoitua porttia SemaphoreSlim.*/?/comma-listat).SignalDispatcher tai AsyncSignalProcessor Käsittelijän sisällä.From Snapshots.cs:
public sealed record EphemeralOperationSnapshot(
long Id,
DateTimeOffset Started,
DateTimeOffset? Completed,
string? Key,
bool IsFaulted,
Exception? Error,
TimeSpan? Duration,
IReadOnlyList<string>? Signals = null,
bool IsPinned = false)
{
public bool HasSignal(string signal) => Signals?.Contains(signal) == true;
}
// For result-capturing coordinators
public sealed record EphemeralOperationSnapshot<TResult>(
long Id,
DateTimeOffset Started,
DateTimeOffset? Completed,
string? Key,
bool IsFaulted,
Exception? Error,
TimeSpan? Duration,
TResult? Result,
bool HasResult,
IReadOnlyList<string>? Signals = null,
bool IsPinned = false);
Tämä on Vain metatiedotHuomaa, mitä on ei tässä:
Vain sen verran, että voin vastata: "Mitä tapahtui, milloin ja toimiko se?" - Ei sen enempää.
.NET antaa useita tapoja tehdä rinnakkaistöitä. Ephemeral-kirjasto vertaa näin:
await Parallel.ForEachAsync(items,
new ParallelOptions { MaxDegreeOfParallelism = 4 },
async (item, ct) => await ProcessAsync(item, ct));
Paras: Yksinkertainen rinnakkaiskäsittely kokoelmista, joissa ei tarvita näkyvyyttä.
Mitä siitä puuttuu:
Käytä Ephemeralia, kun: Tarvitset vianetsintää/tarkkailua, per-avaintilausta tai signaalin reagointia.
var block = new ActionBlock<T>(
async item => await ProcessAsync(item),
new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 4 });
foreach (var item in items)
block.Post(item);
block.Complete();
await block.Completion;
Paras: Monimutkaiset tietovirtaputket, joissa on haaroittumista, yhdistämistä, erittelyä.
Mitä se tekee hyvin:
Käytä TPL-tietovirtaa, kun: Tarvitset monimutkaisia putkistotopologioita (fan-out, fan-in, ehdollinen reititys).
Käytä Ephemeralia, kun: Tarvitset toiminnan seurantaa, yksinkertaisempaa API:tä tai signaalin reagointia.
var channel = Channel.CreateBounded<T>(100);
// Producer
foreach (var item in items)
await channel.Writer.WriteAsync(item);
channel.Writer.Complete();
// Consumer (multiple workers)
var workers = Enumerable.Range(0, 4).Select(async _ =>
{
await foreach (var item in channel.Reader.ReadAllAsync())
await ProcessAsync(item);
});
await Task.WhenAll(workers);
Paras: Tuottaja-kuluttaja-mallit, joissa hallitset molempia osapuolia.
Mitä se tekee hyvin:
Käytä kanavia, kun: Rakennat mukautetun infrastruktuurin ja tarvitset maksimaalista valvontaa.
Käytä Ephemeralia, kun: Haluat toiminnan seurannan ja havainnoitavuuden ilman kattilalevyä.
var policy = Policy
.Handle<HttpRequestException>()
.WaitAndRetryAsync(3, attempt => TimeSpan.FromSeconds(Math.Pow(2, attempt)));
await policy.ExecuteAsync(() => ProcessAsync(item));
Paras: Resilience policy (retriili, virrankatkaisija, aikalisä) yksittäisten toimintojen osalta.
Käytä Pollya, kun: Yksittäisten puheluiden ympärille tarvitaan sietokykyä.
Käytä Ephemeralia, kun: Tarvitset koordinaatiota monissa operaatioissa, joissa on ympäristötietoisuus.
Yhdistä ne: Käytä Pollya Ephemeral-työkehikkosi sisällä per-operaation sietokykyyn.
Paras: Jaettuja viestejä kaikille palveluille kestävin jonoin.
Käytä viestibusseja, kun: Työn on kestettävä prosessin uudelleenkäynnistykset, tarjottava useita palveluja tai vaadittava taattua toimitusta.
Käytä Ephemeralia, kun: Työ on käynnissä, ei tarvitse kestävyyttä, ja haluat kevyttä havainnointia.
Lähestymistapa Rajattu seuranta Signaalit Itsepuhdistava Monimutkaisuus
|----------|:-------:|:--------:|:-------:|:-------:|:-------------:|:----------:|
| Parallel.ForEachAsync . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . .
TPL-tietovirta Korkea-arvoisuus
TV-kanavat Keskipitkät kanavat
Polly, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei, ei
Taustapalvelut Keskipitkällä aikavälillä
MassTranit/NServiceBus Ylhäältä
| Ephemeral Library Alhaalla.
From Rinnakkaiset ephemeral.cs:
// Simple parallel processing with tracking
await items.EphemeralForEachAsync(
async (item, ct) => await ProcessAsync(item, ct),
new EphemeralOptions { MaxConcurrency = 8 });
// With keyed execution (per-user sequential)
await commands.EphemeralForEachAsync(
cmd => cmd.UserId, // Key selector
async (cmd, ct) => await ExecuteCommandAsync(cmd, ct),
new EphemeralOptions
{
MaxConcurrency = 32,
MaxConcurrencyPerKey = 1 // Sequential per user
});
Kuvittele käyttäjäkomentojen käsittely:
Ilman näppäilyä ne voisivat toimia seuraavasti: 1, 4, 2, 5, 3, 6 - välilehdessä.
yy) kanssa. MaxConcurrencyPerKey = 1:
Tämä on per yksikkö peräkkäin, globaalisti rinnakkain - kriittinen järjestelmille, joissa tilauksella on merkitystä kokonaisuuden sisällä.
From EphemeralWorkCoordinator.cs:
await using var coordinator = new EphemeralWorkCoordinator<TranslationRequest>(
async (request, ct) => await TranslateAsync(request, ct),
new EphemeralOptions
{
MaxConcurrency = 8,
MaxTrackedOperations = 500,
EnableDynamicConcurrency = true // Allow runtime adjustment
});
// Enqueue items over time
await coordinator.EnqueueAsync(new TranslationRequest("Hello", "es"));
// Check status anytime
Console.WriteLine($"Pending: {coordinator.PendingCount}");
Console.WriteLine($"Active: {coordinator.ActiveCount}");
// Get snapshots
var snapshot = coordinator.GetSnapshot();
var running = coordinator.GetRunning();
var failed = coordinator.GetFailed();
var completed = coordinator.GetCompleted();
// Control flow
coordinator.Pause(); // Stop pulling new work
coordinator.Resume(); // Continue
// Adjust concurrency at runtime (requires EnableDynamicConcurrency)
coordinator.SetMaxConcurrency(16);
// Pin important operations to survive eviction
coordinator.Pin(operationId);
coordinator.Unpin(operationId);
coordinator.Evict(operationId);
// When done
coordinator.Complete();
await coordinator.DrainAsync();
await using var coordinator = EphemeralWorkCoordinator<Message>.FromAsyncEnumerable(
messageStream, // IAsyncEnumerable<Message>
async (msg, ct) => await ProcessMessageAsync(msg, ct),
new EphemeralOptions { MaxConcurrency = 16 });
await coordinator.DrainAsync();
From EphemeralKeyedWork Coordinator.cs:
await using var coordinator = new EphemeralKeyedWorkCoordinator<string, Command>(
cmd => cmd.UserId, // Key selector
async (cmd, ct) => await ExecuteCommandAsync(cmd, ct),
new EphemeralOptions
{
MaxConcurrency = 32,
MaxConcurrencyPerKey = 1, // Per-user sequential
EnableFairScheduling = true, // Prevent hot user starvation
FairSchedulingThreshold = 10 // Reject if user has 10+ pending
});
// TryEnqueue returns false if fair scheduling rejects
if (!coordinator.TryEnqueue(hotUserCommand))
{
await DeferCommandAsync(hotUserCommand);
}
// Per-key visibility
var pendingForUser = coordinator.GetPendingCountForKey("user-123");
var opsForUser = coordinator.GetSnapshotForKey("user-123");
From EphemeralResultCoordinator.cs:
await using var coordinator = new EphemeralResultCoordinator<SessionInput, SessionResult>(
async (input, ct) =>
{
var fingerprint = await ComputeFingerprintAsync(input.Events, ct);
return new SessionResult(fingerprint, input.Events.Length);
},
new EphemeralOptions { MaxConcurrency = 16 });
await coordinator.EnqueueAsync(session);
coordinator.Complete();
await coordinator.DrainAsync();
// Get just the results (no metadata)
var results = coordinator.GetResults();
// Get snapshots with results + metadata
var snapshots = coordinator.GetSnapshot();
// Get base snapshots without results (privacy-safe)
var baseSnapshots = coordinator.GetBaseSnapshot();
// Filter by success/failure
var successful = coordinator.GetSuccessful();
var failed = coordinator.GetFailed();
From ConcurrencyGates.cs:
Kirjastossa on kaksi rahanvaihtoa valvovaa mekanismia:
SemaphoreSlimQueue<WaiterEntry>UpdateLimit() aika-ajossaEnableDynamicConcurrency = true// Dynamic concurrency adjustment
var coordinator = new EphemeralWorkCoordinator<T>(body,
new EphemeralOptions
{
MaxConcurrency = 4,
EnableDynamicConcurrency = true
});
// Later, based on system load:
coordinator.SetMaxConcurrency(16); // Scale up
coordinator.SetMaxConcurrency(2); // Scale down
From RiippuvuusInjektointi.cs:
Kuten IHttpClientFactoryVoit rekisteröidä nimetyt asetukset:
// Registration
services.AddEphemeralWorkCoordinator<TranslationRequest>("fast",
async (request, ct) => await FastTranslateAsync(request, ct),
new EphemeralOptions { MaxConcurrency = 32 });
services.AddEphemeralWorkCoordinator<TranslationRequest>("accurate",
async (request, ct) => await AccurateTranslateAsync(request, ct),
new EphemeralOptions { MaxConcurrency = 4 });
// Usage
public class TranslationService(IEphemeralCoordinatorFactory<TranslationRequest> factory)
{
private readonly EphemeralWorkCoordinator<TranslationRequest> _fast =
factory.CreateCoordinator("fast");
private readonly EphemeralWorkCoordinator<TranslationRequest> _accurate =
factory.CreateCoordinator("accurate");
}
CreateCoordinator("fast") kahdesti palauttaa saman koordinaattorin"fast" sekä "accurate" Hanki erilliset koordinaattoritKaikki koordinaattorit tarjoavat optimoituja signaalikyselymenetelmiä:
// Get all signals
var signals = coordinator.GetSignals();
// Filter by key (zero-allocation)
var userSignals = coordinator.GetSignalsByKey("user-123");
// Filter by time range
var recentSignals = coordinator.GetSignalsSince(DateTimeOffset.UtcNow.AddMinutes(-5));
var rangeSignals = coordinator.GetSignalsByTimeRange(from, to);
// Filter by signal name or pattern
var rateSignals = coordinator.GetSignalsByName("rate-limit");
var httpSignals = coordinator.GetSignalsByPattern("http.*");
// Check existence (short-circuits on first match)
if (coordinator.HasSignal("rate-limit"))
await ThrottleAsync();
if (coordinator.HasSignalMatching("error.*"))
await AlertAsync();
// Count signals efficiently (no allocation)
var totalSignals = coordinator.CountSignals();
var errorCount = coordinator.CountSignals("error");
var httpCount = coordinator.CountSignalsMatching("http.*");
From EphemeralIdGenerator.cs:
internal static class EphemeralIdGenerator
{
private static long _counter;
private static readonly long _processStart = Environment.TickCount64;
private static readonly int _processId = Environment.ProcessId;
[MethodImpl(MethodImplOptions.AggressiveInlining)]
public static long NextId()
{
var counter = Interlocked.Increment(ref _counter);
// Combine counter with process-unique seed
Span<byte> buffer = stackalloc byte[24];
BitConverter.TryWriteBytes(buffer, _processStart);
BitConverter.TryWriteBytes(buffer.Slice(8), _processId);
BitConverter.TryWriteBytes(buffer.Slice(16), counter);
return unchecked((long)XxHash64.HashToUInt64(buffer));
}
}
stackalloc)Interlocked.Increment)Koordinaattorit eivät varastoi Task referenssejä - vain tissejä:
private int _activeTaskCount;
private readonly TaskCompletionSource _drainTcs;
// In ExecuteItemAsync:
finally
{
// Signal drain when last task completes AND channel iteration is done
if (Interlocked.Decrement(ref _activeTaskCount) == 0 &&
Volatile.Read(ref _channelIterationComplete))
{
_drainTcs.TrySetResult();
}
}
Avainkoordinaattori siivoaa automaattisesti tyhjäkäynnit per avain -semaforit:
private sealed class KeyLock(SemaphoreSlim gate, int maxCount)
{
public SemaphoreSlim Gate { get; } = gate;
public int MaxCount { get; } = maxCount;
public long LastUsedTicks = Environment.TickCount64;
}
// Cleanup runs periodically, removes locks idle > 60 seconds
// Program.cs
var builder = WebApplication.CreateBuilder(args);
// Named coordinators
builder.Services.AddEphemeralWorkCoordinator<TranslationRequest>("fast",
async (req, ct) => await FastTranslateAsync(req, ct),
new EphemeralOptions { MaxConcurrency = 16 });
// Keyed coordinator for per-user commands
builder.Services.AddEphemeralKeyedWorkCoordinator<string, UserCommand>("commands",
cmd => cmd.UserId,
sp =>
{
var handler = sp.GetRequiredService<ICommandHandler>();
return async (cmd, ct) => await handler.HandleAsync(cmd, ct);
},
new EphemeralOptions
{
MaxConcurrency = 32,
MaxConcurrencyPerKey = 1,
EnableFairScheduling = true,
CancelOnSignals = new HashSet<string> { "system-overload" }
});
var app = builder.Build();
// Controller
[ApiController]
[Route("api")]
public class WorkController : ControllerBase
{
private readonly EphemeralWorkCoordinator<TranslationRequest> _translator;
private readonly EphemeralKeyedWorkCoordinator<string, UserCommand> _commands;
public WorkController(
IEphemeralCoordinatorFactory<TranslationRequest> translationFactory,
IEphemeralKeyedCoordinatorFactory<string, UserCommand> commandFactory)
{
_translator = translationFactory.CreateCoordinator("fast");
_commands = commandFactory.CreateCoordinator("commands");
}
[HttpPost("translate")]
public async Task<IActionResult> Translate([FromBody] TranslationRequest request)
{
await _translator.EnqueueAsync(request);
return Ok(new { pending = _translator.PendingCount });
}
[HttpPost("command")]
public IActionResult SubmitCommand([FromBody] UserCommand command)
{
if (!_commands.TryEnqueue(command))
return StatusCode(429, "Too many pending commands for this user");
return Ok();
}
[HttpGet("status")]
public IActionResult GetStatus() => Ok(new
{
translator = new
{
pending = _translator.PendingCount,
active = _translator.ActiveCount,
completed = _translator.TotalCompleted,
failed = _translator.TotalFailed,
hasRateLimit = _translator.HasSignal("rate-limit")
},
commands = new
{
pending = _commands.PendingCount,
active = _commands.ActiveCount,
errorCount = _commands.CountSignalsMatching("error.*")
}
});
}
Olemme rakentaneet täydellisen teloituskirjaston, jossa on:
EphemeralForEachAsync - Yhden laukauksen rinnakkaiskäsittely ja seurantaEphemeralWorkCoordinator - Pitkäikäiset jonotEphemeralKeyedWorkCoordinator - Per-yksikkö peräkkäinen toteutus reilulla aikataulullaEphemeralResultCoordinator - Tulosten vangitseva varianttiIHttpClientFactoryKaava istuu suloiseen kohtaan:
Parallel.ForEachAsyncAmpukaa, älkääkä unohtako.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.