Το **Μέρος πρώτο: Φωτιά και μη Αρκετά. Ξέχνα το.**Εξερευνήσαμε τη θεωρία πίσω από την εφήμερη εκτέλεση - οριοθετημένη, ιδιωτική, αποσφαλματωμένη ροή εργασίας async που θυμούνται αρκετά για να είναι χρήσιμη και στη συνέχεια εξατμίζονται.
Αυτό το άρθρο μετατρέπει αυτό το μοτίβο σε μια επαναχρησιμοποίητη βιβλιοθήκη μπορείτε να πέσει σε οποιοδήποτε έργο .NET.
Η βιβλιοθήκη είναι χωρισμένη σε καλά παραλυμένα αρχεία:
Αρχείο > Σκοπός > > > > Αρχείο > > > > > > > > > > > > > > > > > > > > > > < > > < > > > < > < > > < > > > < > > > > > < > > > > > > < < > > > > > > > > < > > < > > < < > > > > < > > > > > > < < < < > < < < < < < < < < < <
|------|---------|
| EphemeralOptions.cs Διαμόρφωση (νόμισμα, μέγεθος παραθύρου, διάρκεια ζωής, σήματα)
| Ephemeraral Operation.cs □ Εσωτερική παρακολούθηση λειτουργίας με υποστήριξη σήματος
| Στιγμιότυπα. cs Αμετάβλητα αρχεία στιγμιότυπων εκτεθειμένων στους καταναλωτές
| Σήματα. cs Σήμα γεγονότα, διάδοση, περιορισμούς, και το παγκόσμιο SignalSink
| EphemaralIdGenerator.cs Γρήγορη XxHash64-based ID γενεά
| ConcurrencyGates.cs Ο περιορισμός του σταθερού και ρυθμιζόμενου νομίσματος
| StringPatternMatcher.cs Το μοτίβο Glob-style ταιριάζει για το φιλτράρισμα σημάτων
| ΠαράλληλαEphemeral.cs Μέθοδοι στατικής επέκτασης (EphemeralForEachAsync) |
| Ephemeral WorkCoordinator.cs □ Μακροζωία συντονιστής ουράς εργασίας
| EphemeralKeyedWorkCoordinator.cs Ακολουθία εκτέλεσης ανά κλειδί με δίκαιο προγραμματισμό
| EphemeralResultCoordinator.cs Μετάφραση συντονιστή με αποτέλεσμα να συλλαμβάνονται τα αποτελέσματα
| SignalDispatcher.cs Ροή σήματος Async με μοτίβο που ταιριάζει
| Εξαρτήσεις Ένεση. cs Οι μέθοδοι επέκτασης DI και οι εφαρμογές του εργοστασίου
| Παραδείγματα/SignalingHttpClient.cs Το δείγμα εκπομπών σήματος με λεπτό γρανάζιο για τις κλήσεις HTTP
Και ολοκληρωμένες δοκιμές καλύπτει όλες τις ακραίες περιπτώσεις.
Ακούστε τι αντικαθιστούμε:
// ❌ 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...
Και τι φτιάχνουμε:
// ✅ 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();
Ίδια εκτέλεση async. Πλήρης παρατηρητικότητα. Δεν διατηρούνται δεδομένα χρήστη.
Το πιο κοινό μοτίβο - καταχωρήστε έναν συντονιστή σε DI και να το ενέσετε:
// 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│
│ │
└─────────────────────────────────────────────────────────────────┘
Από 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() - χρησιμοποιεί μια προσαρμοσμένη πύλη αντί για SemaphoreSlim.*/?/comma λίστες).SignalDispatcher ή AsyncSignalProcessor Μέσα στον χειριστή.Από Στιγμιότυπα. 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);
Αυτό είναι μόνο μεταδεδομέναΠρόσεξε τι είναι. Όχι, όχι. Εδώ:
Αρκετά για να απαντήσω "τι συνέβη, πότε, και δούλεψε;" - τίποτα περισσότερο.
.NET σας δίνει διάφορους τρόπους για να κάνετε παράλληλη εργασία.
await Parallel.ForEachAsync(items,
new ParallelOptions { MaxDegreeOfParallelism = 4 },
async (item, ct) => await ProcessAsync(item, ct));
Το καλύτερο για: Απλή παράλληλη επεξεργασία συλλογών όπου δεν χρειάζεστε ορατότητα.
Αυτό που της λείπει.:
Χρησιμοποιήστε το Ephemeral όταν: Χρειάζεστε αποσφαλμάτωση/παρατήρηση, ανά κλειδί παραγγελία, ή επεξεργασία αντιδραστική σήματος.
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;
Το καλύτερο για: Συγκρότημα αγωγών ροής δεδομένων με διακλάδωση, συγχώνευση, παρτίδα.
Τι κάνει καλά;:
Χρήση της ροής δεδομένων TPL όταν: Χρειάζεστε πολύπλοκες συγνώμες αγωγών (fan-out, fan-in, condition routing).
Χρησιμοποιήστε το Ephemeral όταν: Χρειάζεστε παρακολούθηση λειτουργίας, απλούστερο API, ή αντιδραστικός συντονισμός σήματος.
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);
Το καλύτερο για: μοτίβα παραγωγού-καταναλωτή όπου ελέγχετε και τις δύο πλευρές.
Τι κάνει καλά;:
Χρήση καναλιών όταν: Χτίζετε προσαρμοσμένη υποδομή και χρειάζεστε μέγιστο έλεγχο.
Χρησιμοποιήστε το Ephemeral όταν: Θέλετε παρακολούθηση λειτουργίας και παρατηρητικότητα χωρίς τον λέβητα.
var policy = Policy
.Handle<HttpRequestException>()
.WaitAndRetryAsync(3, attempt => TimeSpan.FromSeconds(Math.Pow(2, attempt)));
await policy.ExecuteAsync(() => ProcessAsync(item));
Το καλύτερο για: Πολιτικές ανθεκτικότητας (θρησκευτικό, διακόπτης κυκλώματος, χρόνος) για μεμονωμένες λειτουργίες.
Χρησιμοποιήστε την Polly όταν: Χρειάζεστε ανθεκτικότητα γύρω από ατομικές κλήσεις.
Χρησιμοποιήστε το Ephemeral όταν: Χρειάζεστε συντονισμό σε πολλές λειτουργίες με την ευαισθητοποίηση του περιβάλλοντος.
Συνδυάστε τα.: Χρησιμοποιήστε την Polly μέσα στο σώμα εργασίας Ephemeraral σας για την ανθεκτικότητα ανά επιχείρηση.
Το καλύτερο για: Διανεμημένα μηνύματα σε όλες τις υπηρεσίες με ανθεκτικές ουρές.
Χρήση λεωφορείων μηνυμάτων όταν: Η εργασία πρέπει να επιβιώσει επανεκκινεί τη διαδικασία, να καλύπτει πολλαπλές υπηρεσίες, ή να απαιτεί εγγυημένη παράδοση.
Χρησιμοποιήστε το Ephemeral όταν: Η εργασία είναι σε διαδικασία, δεν χρειάζεται αντοχή, και θέλετε ελαφριά παρατηρητικότητα.
Προσεγγίσεις < > > > > > > > > > < > > > > > > > > > > > > > > > > > > > < > > > > < > < > > > < > > < > > > > > < > > > > > > > < > < > > > > > > > > > < > > < > > < > > > < > > > > < > < < < > > < > < < < < < < > > > > > > < > > > > > > > > < < > < < > < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < < <
|----------|:-------:|:--------:|:-------:|:-------:|:-------------:|:----------:|
| Parallel.ForEachAsync ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~
~ TPL Dataflow
Κανάλια
Πόλυ, Ν/Α, Ν/Α, Ν/Α, Ν/Α, Ν/Α, Ν/Α, Ν/Α, Ν/Α.
-Υπηρεσίες υποβάθρου. -Υπηρεσίες υποβάθρου -Μεσαίο
MassTransit/NServiceBus . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . .
| Βιβλιοθήκη Εφήμερων ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~ ~
// 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
});
Φανταστείτε τις εντολές χρηστών επεξεργασίας:
Χωρίς πλήκτρο, αυτά θα μπορούσαν να εκτελεστούν ως: 1, 4, 2, 5, 3, 6 - διαγραμμένα.
Με MaxConcurrencyPerKey = 1:
Αυτό είναι Διαδοχή ανά φορέα, παγκοσμίως παράλληλη - κρίσιμης σημασίας για συστήματα όπου η τάξη έχει σημασία εντός μιας οντότητας.
Από Ephemeral WorkCoordinator.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();
Από EphemeralKeyedWorkCoordinator.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");
Από 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();
Από ConcurrencyGates.cs:
Η βιβλιοθήκη παρέχει δύο μηχανισμούς ελέγχου του νομίσματος:
SemaphoreSlimQueue<WaiterEntry>UpdateLimit() σε χρόνο λειτουργίαςEnableDynamicConcurrency = 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
Από Εξαρτήσεις Ένεση. cs:
Όπως IHttpClientFactory, μπορείτε να καταχωρήσετε τις ρυθμίσεις που ονομάζονται:
// 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") δύο φορές επιστρέφει ο ίδιος συντονιστής"fast" και "accurate" Αποκτήστε ξεχωριστούς συντονιστέςΌλοι οι συντονιστές παρέχουν βελτιστοποιημένες μεθόδους αναζήτησης σημάτων:
// 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.*");
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)Οι συντονιστές δεν αποθηκεύουν Task αναφορές - μόνο μετρητής:
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();
}
}
Ο κλειδωμένος συντονιστής καθαρίζει αυτόματα αδρανείς ημιαγωγούς ανά κλειδί:
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.*")
}
});
}
Φτιάξαμε μια πλήρη εφήμερη βιβλιοθήκη εκτέλεσης με:
EphemeralForEachAsync - Μία-shot παράλληλη επεξεργασία με ανίχνευσηEphemeralWorkCoordinator - Μακροζωείς ορατές ουρέςEphemeralKeyedWorkCoordinator - Διακεκριμένη εκτέλεση ανά φορέα με δίκαιο προγραμματισμόEphemeralResultCoordinator - Παραλλαγή σύλληψης αποτελεσμάτωνIHttpClientFactoryΤο μοτίβο κάθεται σε ένα γλυκό σημείο:
Parallel.ForEachAsyncΦωτιά... και μην ξεχνάς.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.