में **पार्ट 1: आग और मत काफी भूल जाना**हमने एक - एक करके इस सिद्धांत की जाँच की, जो निजी तौर पर, डीबगी का काम करता है ।
इस लेख से पता चलता है कि आप किसी भी तरह की लाइब्रेरी में जा सकते हैं ।
अब यह विशेष रूप से अल्पतम पैकेज में है 20 से अधिक।.
लाइब्रेरी को अच्छी तरह से जमा किया गया फ़ाइलों में विभाजित किया गया है:
IMTA फ़ाइल BAR
|------|---------|
| एपीटील विकल्प. सेंटीमीटर कॉन्फ़िगरेशन (अनुप्रयोग, विंडो आकार, जीवन, संकेत)
| एपीड्रॉडिशन. SICAR आंतरिक ऑपरेशन संकेत समर्थन के साथ ट्रैक किया जा रहा है
| स्नेपशॉट्स.css फूजनिक स्नैपशॉट रिकॉर्डों को उपभोक्ताओं के संपर्क में लाया जा सकता है
| सिग्नल.cs सिग्नल घटना, फूटना, प्रतिबन्ध, और वैश्विक सिग्नल- ओएस
| एपीडैम्टर सेंटीमीटर तेज XxHM64- आधारित IDEALLLANK
| दोष लगाना । डीवीडी+आर( डबल्यू) स्थिर तथा समायोजित किया जा सकता है
| स्ट्रिंग इंटरनेशनलमैटर. संकेत फिल्टरिंग के लिए सेल- दर- दर- दर- दर- रंग पैटर्न
| समानांतर प्रेसमेराल. स्थैतिक एक्सटेंशन विधियाँ (i)EphemeralForEachAsync) |
| एपम्रॉमवर्क कोर. बहुतेरे लंबे समय के कार्य कतार कोऑर्डिनेटर
| एपीट- कम्पेटर PRECT पर- कुंजी अनुक्रम उचित सारिणी के साथ चलाना
| एपमर्मल (Veteowordordin.s) INVEGES परिणाम- सेटअप किए जा रहे संयोजक
| सिग्नल अमेरिका.cs Smughughue सिग्नल पैटर्न से मेल खाने के साथ
| डिपेंडेंसी इन्जेक्शन. INVES विस्तार विधि और फैक्टरी
| उदाहरण/ चिह्न एचटीटीपी कॉल के लिए बीजीय नमूना प्लगिन
और विस्तार जाँच सभी किनारे के मामलों को कवर.
यहाँ हम क्या बदल रहे हैं:
// ❌ 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();
एक ही बार में चलाने के लिए. पूर्ण. कोई उपयोक्ता डाटा नहीं बचा.
सबसे सामान्य पैटर्न - 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│
│ │
└─────────────────────────────────────────────────────────────────┘
से एपीटील विकल्प.:
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.*/?/कमा सूची.SignalDispatcher या AsyncSignalProcessor हैंडलर के अंदर.से स्नेपशॉट्स.css:
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);
यह है सिर्फ मेटाडाटाध्यान दें कि क्या है नहीं यहाँ:
बस जवाब देने के लिए पर्याप्त, जब, और यह काम किया था? - और कुछ भी नहीं.
.netT आप समानांतर काम करने के लिए कई तरीके देता है। यहाँ कैसे एपिडिडल लाइब्रेरी तुलना करता है:
await Parallel.ForEachAsync(items,
new ParallelOptions { MaxDegreeOfParallelism = 4 },
async (item, ct) => await ProcessAsync(item, ct));
के लिए उत्तम: संग्रह की सरल समानांतर प्रक्रिया जहाँ आप की दृश्यता की जरूरत नहीं है.
इसमें कोई कमी नहीं:
जब एपेक्ट्डल इस्तेमाल करें तो इस्तेमाल करें: आपको डिबगिंग/ ब्लिकमेंट की आवश्यकता है, प्रति-key आदेश, या संकेत-सक्रिय प्रक्रिया की जरूरत है.
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;
के लिए उत्तम: साफ - सफाई के साथ - साथ जटिल जानकारी भी आती है ।
यह अच्छी तरह से क्या करता है:
जब टी. एल. एल.): आपको जटिल इंफेक्शन (f- आउट, प्रशंसक-in) की जरूरत है.
जब एपेक्ट्डल इस्तेमाल करें तो इस्तेमाल करें: आपको ट्रैक की जरूरत है, आसान एपीआई, या संकेत एम्बिलेशन की जरूरत है.
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);
के लिए उत्तम: निर्माता-कंकर पैटर्न जहां आप दोनों पक्षों को नियंत्रित करते हैं।
यह अच्छी तरह से क्या करता है:
जब चैनल्स इस्तेमाल करें: तुम इन्फ्रास्फीति का रिवाज़ बना रहे हो और अधिकतम नियंत्रण की जरूरत है.
जब एपेक्ट्डल इस्तेमाल करें तो इस्तेमाल करें: आप बिना साबुन के ट्रैकिंग और ब्लाइडमेंट करना चाहते हैं.
var policy = Policy
.Handle<HttpRequestException>()
.WaitAndRetryAsync(3, attempt => TimeSpan.FromSeconds(Math.Pow(2, attempt)));
await policy.ExecuteAsync(() => ProcessAsync(item));
के लिए उत्तम: अलग - अलग ऑपरेशनों के लिए रेस्टिंग पॉलिसी (अवरी, सर्किट ब्रेकर, टाइम) ।
जब पोल इस्तेमाल करें: आपको हर व्यक्ति से बात करने की ज़रूरत है ।
जब एपेक्ट्डल इस्तेमाल करें तो इस्तेमाल करें: आपको बहुत - से ऑपरेशनों में सावधानी बरतने की ज़रूरत है ।
उन्हें सम्मिलित करें: हर व्यक्ति के लिए अपने एपिडीटीम काम बॉडी के अंदर पोली से काम लीजिए.
के लिए उत्तम: बहुत - से देशों में साक्षियों के काम पर पाबंदी लगी हुई थी ।
जब संदेश विंडो इस्तेमाल करें: काम को फिर से शुरू करने की प्रक्रिया से, बहुत सी सेवाएँ शुरू करने की ज़रूरत है या फिर उसे पूरा करने की ज़रूरत है ।
जब एपेक्ट्डल इस्तेमाल करें तो इस्तेमाल करें: काम प्रक्रिया में है, scrrerircenty की आवश्यकता नहीं है, और आप हल्का bocrscrervivy चाहते हैं.
MRITCACT CONT TECKT TANT T-CKT TECKT TECT TECT TECKT TECT TKT TECT TENT T-key SECKS SECKS SECKSTCKS SICKRICKS SICKRICKCKS TANT TENT TANT TENT TENT TENT TENT TENT TENT TE(KCKCKCKCKSECKCKCKCKCKCKCKCKCKSTCKSTCK(K)
|----------|:-------:|:--------:|:-------:|:-------:|:-------------:|:----------:|
| Parallel.ForEachAsync ख़ीनएल © 2002- 2006 घड़ियात लिखा आपसी सहयोग
बेलारूसी डाटा स्थाप इससे सरियाना लिखा हुआ है ।
बख्शिश चैनल अनक़रीब ही निष्क्रिय कर देगा
बिलिक पोलली निया/ALLLYEAL N/ALLLLYELLL N/ALLLLLL NALLLYYELLLLL NAR_ /ALLLLLLLLLLLLLL N/ALLLLLLLLLLLL
बुरी तरह से अधिक पृष्ठभूमि सेवा
राहित द्रव्यमान/ NELLLALLLLLLLYALLLLLLLLLYY_BAR_
| एपाइडल लाइब्रेरी बुर्किना फासो घड़ू आदमियों की आपस में आदमियों की मौत हो चुकी है । ( g04 7 / 22)
// 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
});
उपयोक्ता कमांड्स कल्पना करें:
कुंजी के बिना, ये शायद इस प्रकार कार्य करें: १, ४, ५, ३, ३, ६ - ४.
के साथ MaxConcurrencyPerKey = 1:
यह है प्रति- केंद्र अनुक्रम, विश्वव्यापी समानांतर - सिस्टमों के लिए महत्वपूर्ण जहां किसी एंटिटी के भीतर व्यवस्था व्यवस्था.
से एपम्रॉमवर्क कोर.:
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();
से एपीट- कम्पेटर:
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");
से एपमर्मल (Veteowordordin.s):
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();
से दोष लगाना ।:
लाइब्रेरी दो सुग्राहक नियंत्रण यांत्रिकी प्रदान करता है:
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
जैसे 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 संदर्भ - सिर्फ काउंटर: (n)
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();
}
}
कुंजी संयोजक स्वतः निष्क्रियता को साफ करता है प्रति-key रीस्फोर्स:
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 - ट्रैक के साथ एक स्नेपशॉट प्रक्रियाEphemeralWorkCoordinator - लंबे समय का दृश्य कतारEphemeralKeyedWorkCoordinator - सही अनुसूचित के साथ PRATY अनुक्रमEphemeralResultCoordinator - परिणाम-Sming चरIHttpClientFactoryपैटर्न एक प्यारी जगह पर बैठता है:
Parallel.ForEachAsyncआग ... और काफी भूल नहीं है.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.