داخل الجزء 1: النار ولا الـ مـا مـا قبل انسي النسيان،استكشفنا النظرية وراء التنفيذ الإيفيميري - مُحَدَّد، خاص، مُطَلَّب async سير العمل الذي يَتذكّرُ فقط بما فيه الكفاية لِكي يَكُونَ مفيدَ وبعد ذلك يتبخر.
هذه المادة تحول هذا النمط إلى مكتبة قابلة لإعادة الاستخدام يمكنك أن تسقط في أي مشروع NET.
هذا هو الآن في معظم الغالب liflucyd.ephimmerals Nuget حزمة أيضا أكثر من 20.
وتنقسم المكتبة إلى ملفات مصممة بشكل جيد:
□ الملف / الغرض من
|------|---------|
| أولاً- اجتماع الدول الأطراف في o إعدادات (عملة، حجم نافذة، عمر، إشارات)
| التنفيذية التنفيذية التنفيذية التنفيذية التنفيذية □ تتبع العملية الداخلية مع دعم إشارة
| مناقصات. سجلات سريعة غير ثابتة تعرض للمستهلكين
| إشارات الوصلات • الأحداث الاشاراتية والدعاية والمعوقات والوصلة الاشاراتية العالمية
| مُجِدِّل. ccs توليد بطاقة تعريف ذاتية تعتمد على XxHash64
| الأطراف(أ)(ج)(ج) ● معامِل ثابت وقابل للتعديل مقيِد
| سلسلة سِطر سِطر P P P P Matn Mattchr.cs □ مطابقة نمط نمط أسلوب غلوب لرشح الإشارة □ □ □ □ □ □ □ □ □
| )ج( □ طرق التوسّعEphemeralForEachAsync) |
| منسق الأمم المتحدة المنسق)ج( منسق لمناصب العمل الطويلة الأجل
| منسق أعمال الأمم المتحدة التنفيذ بالتتابع على أساس المفتاح مع جدولة عادلة
| منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق (ب) متغيرات المنسقين الذين يحصلون على النتائج(أ)
| اشارة الاشارة الدالة Asyync إشارة متجهة مع نمط مطابقة
| المحتويات (تابع) ● طرق توسيع نطاق نظام DDI وعمليات تنفيذ المصانع
| الأمثلة/الأمثلة/الأمثلة المتسلسلـة: HttttpClient.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 تنفيذ كامل مُلاحظة لا مستخدم البيانات مُستبقى.
أكثر النمط شيوعاً - تسجيل منسق في نظام المعلومات التصميمية وحقنه:
// 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 -العامل الداخلىمن طراز من طراز مناقصات.:
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));
أفضل لـ: معالجة متوازية بسيطة للمجموعات حيث لا تحتاج إلى رؤية.
ما ينقصه:
استخدم متى: تحتاج إلى التصحيح/الملاحظة، أو طلب المفتاح الواحد، أو معالجة الإشارات-التفاعلية.
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;
أفضل لـ: خطوط أنابيب تدفق البيانات المعقدة مع التفرع، والدمج، والدفع.
ما تقوم به جيداً:
****: تحتاج إلى طوبولوجيات معقدة في خطوط الأنابيب (الانفان-ت، مروحة في، توجيه مشروطة).
استخدم متى: تحتاج إلى تتبع العمليات، أو أبسط 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);
أفضل لـأنماط الاستهلاك المنتج حيث تسيطر على كلا الجانبين.
ما تقوم به جيداً:
****أنت تبني البنية التحتية الجمركية وتحتاج إلى الحد الأقصى من التحكم.
استخدم متىتريد تتبع العمليات و المراقبة بدون طاولة الغلاية
var policy = Policy
.Handle<HttpRequestException>()
.WaitAndRetryAsync(3, attempt => TimeSpan.FromSeconds(Math.Pow(2, attempt)));
await policy.ExecuteAsync(() => ProcessAsync(item));
أفضل لـ: السياسات المتعلقة بالقدرة على المنافسة (المراجعة، ومفرق الدوائر، والمهلة الزمنية) لكل عملية على حدة.
****: تحتاج إلى المرونة حول المكالمات الفردية.
استخدم متى: تحتاج إلى التنسيق عبر العديد من العمليات مع الوعي المحيط.
يُمَزّجُهم: استخدام بولي داخل هيئة العمل الإكهوارية الخاصة بك لمقاومة كل عملية.
أفضل لـ: توزيع الرسائل عبر الخدمات مع طابور دائم.
استخدام رسالة الرسالة متى: العمل يجب أن يبقى على قيد الحياة عمليات استئناف التشغيل، أو يشمل خدمات متعددة، أو يتطلب تنفيذ مضمون.
استخدم متى: العمل هو في عملية التشغيل، لا يحتاج إلى متانة، وأنت تريد مراقبة خفيفة الوزن.
□ النهج المتبع □ التتبع المقيد للإشارات التي يشير إليها كل مركز رئيسي □ معقَّد مُنْقِس ذاتي النُقْص
|----------|:-------:|:--------:|:-------:|:-------:|:-------------:|:----------:|
| Parallel.ForEachAsync لا □ لا □ لا □ لا □ لا □ لا □ لا □ لا □ لا □ لا □ لا □
□ تدفق البيانات في حالة عدم وجود بيانات
□ قنان □ □ □ □ □ □ □
لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا لا
خدمات معلومات أساسية □ □ □ □ □ □ □ □ □ □ □ □ □ متوسطة □ □ □ □ □
□ ماس عبرت/نرفيـع البوابـز □ ماس عبرت/نافـسـفـيـبـوبـوس □ □ ماسـت عبرت/نافـسـفـيـبـوس
| المكتبة الصفحة o o o o o o o lo o
من طراز من طراز )ج(:
// 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:
هذا التسلسلي، المتوازي على الصعيد العالمي - ذات أهمية حيوية للنظم التي يكون فيها الطلب ذا شأن داخل كيان ما.
من طراز من طراز منسق الأمم المتحدة المنسق)ج(:
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");
من طراز من طراز منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق/منسق:
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.*");
من طراز من طراز مُجِدِّل. ccs:
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 - معالجة واحدة لرصاصة واحدة مع التتبعEphemeralWorkCoordinator - الطاولات الطويلة الأمد التي تُثقَرEphemeralKeyedWorkCoordinator - التنفيذ بالتعاقب التبعية مع جدولة جدولة عادلةEphemeralResultCoordinator - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - -IHttpClientFactoryالنمط يَجْلسُ في a بقعة حلوة:
Parallel.ForEachAsyncالنار... ولا تنساها تماماً.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.