Skip to content

บทที่ 8 — Background Service + Channels + Hosted Service

← บทที่ 7 | สารบัญ

ASP.NET Core มี IHostedService ซึ่งเป็น interface (ข้อตกลงรูปแบบ method) กลางสำหรับงานเบื้องหลัง — สิ่งที่ implement interface นี้จะถูก host รันให้อัตโนมัติตลอดอายุแอป

ตัวอย่าง use case:

  • cron-like — ทำงานตามเวลาเป็นรอบ (เช่น clean log ทุกชั่วโมง)
  • queue worker — ตัวประมวลผลงานจากคิวภายในแอป
  • message consumer — ตัวรับ message จากระบบคิวภายนอก (Kafka, RabbitMQ)
  • real-time stream — รับ stream ข้อมูลต่อเนื่อง

🚀 โซนขั้นสูง — มือใหม่ข้ามได้ บทนี้เป็นเรื่องงานเบื้องหลัง (queue, message consumer, cron, distributed lock) ส่วนต้น ๆ (หัวข้อ 1-4: IHostedService, BackgroundService, timer, Channel) เป็นพื้นฐานที่ควรรู้ ส่วนหลัง (Kafka/RabbitMQ, distributed lock, outbox) ค่อยกลับมาตอนทำระบบที่ scale หลายเครื่องหรือต่อกับ message queue จริง


1. IHostedService — รากของทุก background

csharp
public class HelloHostedService : IHostedService
{
    public Task StartAsync(CancellationToken ct)
    {
        Console.WriteLine("started");
        return Task.CompletedTask;
    }

    public Task StopAsync(CancellationToken ct)
    {
        Console.WriteLine("stopping");
        return Task.CompletedTask;
    }
}

builder.Services.AddHostedService<HelloHostedService>();

Lifecycle:

  1. StartAsync รันก่อน app start listen
  2. App ทำงาน
  3. SIGTERM → StopAsync + กรอบ HostOptions.ShutdownTimeout

⚠️ StartAsync ไม่ควร block — ถ้าทำงานนาน → app ขึ้นช้า/ไม่ขึ้น — spawn Task แล้ว return

⚠️ ถ้า StartAsync ของ IHostedService ใด throw exception หลุดออกมา host จะ fail to start ทั้งหมด ใน .NET 6+ — ดังนั้นห่อ initialization logic ด้วย try-catch เสมอ และให้ fail fast พร้อม error message ชัดเจน

สำหรับ initialization ที่ใช้เวลานาน แยกอีกประเด็นหนึ่ง — ควรใช้ IHostApplicationLifetime.ApplicationStarted เพื่อดีเฟอร์ (เลื่อน) งานไว้จนหลัง host พร้อมแล้ว จะได้ไม่ไปบล็อกช่วง startup


2. BackgroundService (preferred)

BackgroundService = abstract class ที่ wrap IHostedService ให้ง่ายขึ้น:

💡 abstract class คือ class ที่สร้าง instance ตรง ๆ ไม่ได้ — ต้อง inherit แล้ว implement method ที่ยังไม่มี (เช่น ExecuteAsync) ก่อน ดังนั้นจึงเขียน new BackgroundService() ตรง ๆ ไม่ได้

csharp
public class CleanupService : BackgroundService
{
    private readonly ILogger<CleanupService> _logger;
    private readonly IServiceScopeFactory _scopeFactory;   // ใช้ scope factory ไม่ใช่ root provider

    public CleanupService(ILogger<CleanupService> logger, IServiceScopeFactory scopeFactory)
    {
        _logger = logger;
        _scopeFactory = scopeFactory;
    }

    protected override async Task ExecuteAsync(CancellationToken ct)
    {
        using var timer = new PeriodicTimer(TimeSpan.FromHours(1));
        try
        {
            while (await timer.WaitForNextTickAsync(ct))
            {
                try
                {
                    await using var scope = _scopeFactory.CreateAsyncScope();   // async scope รองรับ IAsyncDisposable (เช่น DbContext)
                    var db = scope.ServiceProvider.GetRequiredService<AppDbContext>();
                    var deleted = await db.Logs
                        .Where(l => l.CreatedAt < DateTime.UtcNow.AddDays(-30))
                        .ExecuteDeleteAsync(ct);   // ส่ง ct ต่อ — bulk delete ยกเลิกได้
                    _logger.LogInformation("Deleted {N} old logs", deleted);
                }
                catch (OperationCanceledException) when (ct.IsCancellationRequested)
                {
                    throw;   // ปล่อย propagate ออกไป shutdown loop
                }
                catch (Exception ex)
                {
                    _logger.LogError(ex, "cleanup failed");
                }
            }
        }
        catch (OperationCanceledException) when (ct.IsCancellationRequested)
        {
            // shutdown ปกติ — ไม่ใช่ error
            _logger.LogInformation("CleanupService stopping gracefully");
        }
    }

    public override async Task StopAsync(CancellationToken ct)
    {
        _logger.LogInformation("CleanupService received stop signal");
        // ถ้ามี state ที่ต้อง flush ทำตรงนี้ก่อนเรียก base
        await base.StopAsync(ct);
    }
}

builder.Services.AddHostedService<CleanupService>();

Key points:

  • ExecuteAsync รันเป็น Task แบบขนานกับ app — ไม่ block startup
  • รับ CancellationToken — ต้อง honor (เคารพ/ทำตาม) สำหรับ graceful shutdown
  • inject IServiceScopeFactory แล้วเรียก CreateAsyncScope() ถ้าต้องการ scoped service (เช่น DbContext)
  • ส่ง ct ต่อให้ async call ทุกชั้น รวมถึง ExecuteDeleteAsync(ct) — ไม่งั้น bulk operation ยกเลิกไม่ได้
  • catch ทุก exception ที่ไม่ใช่ OperationCanceledException — ไม่งั้น service หยุดทำงาน
  • override StopAsync ถ้าต้องทำ cleanup เพิ่ม (เช่น flush buffer, ปิด connection)

⚠️ BackgroundServiceExceptionBehavior.Ignore — การตั้งค่านี้หมายความว่า service ที่ crash จะ หยุดทำงานเงียบ ๆ โดยไม่มี recovery อัตโนมัติ host ยังอยู่ แต่ background task ไม่ทำงานอีกต่อไป

นี่อาจเลวร้ายกว่าการ crash ตรง ๆ เสียอีก เพราะเป็น silent failure (พังแบบไม่มีใครรู้) — ไม่มี error, ไม่มี alert, ระบบดูเหมือนปกติทั้งที่ background ตายไปแล้ว

แนวทางที่ดีกว่าคือ catch exception ทุกตัวใน ExecuteAsync loop แล้ว implement retry/recovery logic เอง แทนที่จะอาศัย Ignore


3. PeriodicTimer — Cron Replacement

csharp
using var timer = new PeriodicTimer(TimeSpan.FromMinutes(5));
while (await timer.WaitForNextTickAsync(ct))
{
    await DoWork(ct);
}

ดีกว่า Task.Delay ใน loop — ไม่ drift, cancel ผ่าน ct, async-friendly

💡 ต้องการ cron expression จริง? ใช้ Quartz.NET หรือ Hangfire


4. Queue Worker — Channel<T>

Channel<T> ใน System.Threading.Channels = ท่อส่งงาน throughput สูง (ปริมาณงานต่อเวลา) ระหว่าง 2 ฝั่ง:

💡 System.Threading.Channels มาใน .NET base library แล้ว — ไม่ต้อง dotnet add package เพิ่ม

  • producer (ตัวผลิตงาน) — เช่น endpoint ที่รับ request แล้วโยนงานเข้าคิว
  • consumer (ตัวประมวลผลงาน) — background service ที่อ่านงานออกมาทำ
csharp
// register channel
builder.Services.AddSingleton(_ => Channel.CreateBounded<EmailJob>(new BoundedChannelOptions(1000)   // 1000 = ความจุสูงสุดของคิว (capacity)
{
    FullMode = BoundedChannelFullMode.Wait,
}));
builder.Services.AddSingleton<IEmailQueue, EmailQueue>();
builder.Services.AddHostedService<EmailWorker>();
csharp
public record EmailJob(string To, string Subject, string Body);

public interface IEmailQueue
{
    ValueTask Enqueue(EmailJob job, CancellationToken ct);
}

public class EmailQueue : IEmailQueue
{
    private readonly Channel<EmailJob> _ch;
    public EmailQueue(Channel<EmailJob> ch) => _ch = ch;
    public ValueTask Enqueue(EmailJob job, CancellationToken ct) =>
        _ch.Writer.WriteAsync(job, ct);
}

public class EmailWorker : BackgroundService
{
    private readonly Channel<EmailJob> _ch;
    private readonly IServiceScopeFactory _scopeFactory;
    private readonly ILogger<EmailWorker> _logger;

    public EmailWorker(Channel<EmailJob> ch, IServiceScopeFactory scopeFactory, ILogger<EmailWorker> logger)
    {
        _ch = ch; _scopeFactory = scopeFactory; _logger = logger;
    }

    protected override async Task ExecuteAsync(CancellationToken ct)
    {
        await foreach (var job in _ch.Reader.ReadAllAsync(ct))
        {
            try
            {
                await using var scope = _scopeFactory.CreateAsyncScope();
                var sender = scope.ServiceProvider.GetRequiredService<IEmailSender>();
                await sender.Send(job.To, job.Subject, job.Body, ct);
            }
            catch (OperationCanceledException) when (ct.IsCancellationRequested)
            {
                throw;   // shutdown — drain loop
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "send failed: {To}", job.To);
                // ตรงนี้ใส่ logic retry หรือส่งเข้า DLQ (ดูคำอธิบายใต้ block)
            }
        }
    }
}

💡 DLQ (Dead Letter Queue) = คิวสำรองสำหรับเก็บงานที่ทำไม่สำเร็จซ้ำ ๆ ไว้ตรวจทีหลัง แทนที่จะทิ้งหายเงียบ หรือ retry ไม่หยุดจนชนกับงานปกติ — โดยปกติจะมีคนมา monitor DLQ และทำ retry หรือ alert ด้วยมือในภายหลัง

ใน controller:

csharp
public async Task Notify(IEmailQueue queue, EmailJob job, CancellationToken ct)
{
    await queue.Enqueue(job, ct);
    // response เร็ว — email ส่งจริงใน background
}

Bounded vs Unbounded:

  • Bounded (CreateBounded) — จำกัดขนาดคิว, มี backpressure (แรงดันย้อนกลับ — เมื่อคิวเต็ม จะหน่วง/บล็อก producer ไม่ให้ผลิตงานเร็วเกินกว่าที่ consumer รับไหว)
  • Unbounded — ไม่จำกัดขนาด (อันตราย ถ้า producer เร็วกว่า consumer คิวจะโตไม่หยุดจน memory เต็ม)

FullMode (พฤติกรรมเมื่อคิวเต็ม):

  • Wait — producer รอ (block) จนมี slot ว่าง
  • DropOldest/DropNewest/DropWrite — ทิ้งงาน (discard) ตัวเก่าสุด/ใหม่สุด/ตัวที่กำลังเขียน

5. Parallel Worker

ถ้า worker ตัวเดียวประมวลผลไม่ทัน job ที่เข้ามา → เพิ่ม worker ขนานได้ โดยรันหลาย task ที่อ่านจาก channel เดียวกัน

จุดสำคัญ:

  • framework กระจาย job ให้ worker ตัวที่ว่างก่อน
  • throughput (ปริมาณงานต่อเวลา) สูงขึ้นตามจำนวน worker
  • ระวัง shared state (ข้อมูลที่ใช้ร่วมกันระหว่าง worker) — ต้อง thread-safe (ปลอดภัยเมื่อหลาย thread เข้าถึงพร้อมกัน) ไม่งั้นข้อมูลพังเงียบ
csharp
protected override async Task ExecuteAsync(CancellationToken ct)
{
    var tasks = Enumerable.Range(0, 4).Select(_ => Task.Run(async () =>
    {
        await foreach (var job in _ch.Reader.ReadAllAsync(ct))
            await Process(job, ct);
    }, ct));
    await Task.WhenAll(tasks);
}

4 worker concurrent อ่านจาก channel เดียวกัน (Channel<T> thread-safe โดยธรรมชาติ)

💡 ct ที่ส่งให้ Task.Run ทำหน้าที่ cancel task ก่อนเริ่ม เท่านั้น

การยกเลิกงานระหว่างรันจริง ๆ ต้องทำผ่าน ct ที่ส่งต่อให้ ReadAllAsync(ct) ข้างในแต่ละ task

ถ้า loop ข้างในตัวใดตัวหนึ่ง throw exception หลุดออกมา ข้อผิดพลาดจะถูกเก็บไว้ใน Task ตัวนั้น แล้ว propagate ออกมาผ่าน Task.WhenAll — ดังนั้นควร catch exception ในลูปเองด้วย ไม่ปล่อยให้หลุดออกมาเฉย ๆ


6. Kafka / RabbitMQ Consumer

use case ที่พบบ่อยที่สุดของ background service คือทำตัวเป็น message consumer

โครงสร้าง:

  • host BackgroundService ที่ subscribe (สมัครรับ) queue/topic (ช่องทางส่ง message ในระบบคิว)
  • process message ที่เข้ามาตลอดอายุแอป
  • ปิด connection ให้สะอาดตอน shutdown

จุดที่พลาดง่าย — manual ack/commit:

  • ack/commit = การยืนยันกับระบบคิวว่ารับ/ทำ message สำเร็จแล้ว
  • manual = ยืนยันเองเมื่อทำงานเสร็จจริง ไม่ใช่ยืนยันอัตโนมัติตอนรับ message
  • ถ้าตั้ง auto-ack แล้วพังกลางทาง → message หายไม่มีคนทำ
  • manual ack ปลอดภัยกว่า — ถ้าพัง message ยังอยู่ในคิว ถูกส่งให้ consumer อื่นทำต่อ

RabbitMQ

bash
dotnet add package RabbitMQ.Client

⚠️ ตัวอย่างนี้ใช้ API ของ RabbitMQ.Client v7+ (IChannel, CreateChannelAsync, AsyncEventingBasicConsumer.ReceivedAsync) — ตรวจ version ใน .csproj ก่อน copy โค้ด: <PackageReference Include="RabbitMQ.Client" Version="x.x.x" />

Versionใช้ API
6 และก่อนหน้าIModel + CreateModel
7+IChannel + CreateChannelAsync
csharp
using System.Text;
using System.Text.Json;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;

public class RabbitConsumer : BackgroundService
{
    private readonly ILogger<RabbitConsumer> _logger;
    private IConnection? _conn;
    private IChannel? _ch;

    public RabbitConsumer(ILogger<RabbitConsumer> logger) => _logger = logger;

    protected override async Task ExecuteAsync(CancellationToken ct)
    {
        // ⚠️ ตัวอย่างนี้ hardcode HostName = "rabbit" เพื่อความสั้น
        // ของจริงควร inject IConfiguration แล้วอ่านจาก appsettings.json
        // เช่น: var factory = new ConnectionFactory { HostName = _config["RabbitMQ:Host"] };
        var factory = new ConnectionFactory { HostName = "rabbit" };
        _conn = await factory.CreateConnectionAsync(ct);
        _ch = await _conn.CreateChannelAsync(cancellationToken: ct);
        await _ch.QueueDeclareAsync("emails", durable: true, exclusive: false, autoDelete: false, cancellationToken: ct);

        var consumer = new AsyncEventingBasicConsumer(_ch);
        consumer.ReceivedAsync += async (_, ea) =>
        {
            if (_ch is null) return;   // guard: ป้องกัน null dereference ถ้า channel ยังไม่พร้อม
            try
            {
                var body = Encoding.UTF8.GetString(ea.Body.Span);
                var job = JsonSerializer.Deserialize<EmailJob>(body);
                if (job is null)   // JSON literal "null" หรือ payload ผิดรูปแบบ deserialize ได้ null โดยไม่ throw
                {
                    _logger.LogWarning("deserialize ได้ null, nack ทิ้ง");
                    await _ch.BasicNackAsync(ea.DeliveryTag, false, requeue: false);
                    return;
                }
                // process... (ของจริงควรโยนเข้า Channel<T> ภายในแล้วให้ worker pool ประมวลผล
                // เพื่อไม่บล็อก dispatcher ของ AsyncEventingBasicConsumer)
                await _ch.BasicAckAsync(ea.DeliveryTag, false);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "rabbit consume failed, requeue");
                await _ch.BasicNackAsync(ea.DeliveryTag, false, requeue: true);
            }
        };

        await _ch.BasicConsumeAsync("emails", autoAck: false, consumer, ct);

        // keep alive จนกว่า ct จะ cancel — TaskCompletionSource ที่ผูกกับ ct สะอาดกว่า Task.Delay infinite
        // TaskCompletionSource = ตัวสร้าง Task ที่เราสั่งให้ "เสร็จ" (complete) เองด้วยมือ ไม่ใช่รอ async method คืนค่า
        // ใช้ตรงนี้เพื่อให้ ExecuteAsync ค้างรอเฉย ๆ จนกว่า ct จะถูกยกเลิก แล้วค่อยปล่อยให้ทำงานต่อ
        var tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
        using var reg = ct.Register(() => tcs.TrySetResult());
        await tcs.Task;
    }

    public override async Task StopAsync(CancellationToken ct)
    {
        _logger.LogInformation("RabbitConsumer stopping, closing channel");
        // pattern นี้ใช้ nullable field (_conn/_ch) ที่ set ใน ExecuteAsync
        // เหตุผลที่ต้อง nullable: ถ้า ExecuteAsync throw ก่อนจะ assign ค่าสำเร็จ StopAsync ก็ยังทำงานได้ปลอดภัยผ่าน null check ด้านล่าง
        // แนวทางที่ clean กว่านี้คือย้าย cleanup ไปไว้ใน try-finally ของ ExecuteAsync เอง แทนที่จะพึ่ง nullable field
        try
        {
            if (_ch is not null) await _ch.CloseAsync();
            if (_conn is not null) await _conn.CloseAsync();
        }
        catch (Exception ex)
        {
            _logger.LogWarning(ex, "error closing rabbit connection");
        }
        await base.StopAsync(ct);
    }
}

Kafka

bash
dotnet add package Confluent.Kafka
csharp
public class KafkaConsumer : BackgroundService
{
    private readonly ILogger<KafkaConsumer> _logger;
    public KafkaConsumer(ILogger<KafkaConsumer> logger) => _logger = logger;

    protected override Task ExecuteAsync(CancellationToken ct)
    {
        // consumer.Consume(ct) เป็น blocking sync call ของ Confluent.Kafka
        // ห้ามเรียกตรงใน async method — จะ block thread pool และทำให้ startup ไม่ complete
        // → ห่อด้วย Task.Run ให้รันใน dedicated thread แทน
        // Task.Run คืน Task ที่แทน background work — BackgroundService await Task นี้
        // ถ้า ConsumeLoop throw exception มันจะ propagate ผ่าน Task และ trigger
        // BackgroundServiceExceptionBehavior (StopHost by default ใน .NET 6+)
        // → ให้ catch exception ใน ConsumeLoop เองให้ครบ
        return Task.Run(() => ConsumeLoop(ct), ct);
    }

    private async Task ConsumeLoop(CancellationToken ct)
    {
        var cfg = new ConsumerConfig
        {
            BootstrapServers = "kafka:9092",
            GroupId = "users-service",
            AutoOffsetReset = AutoOffsetReset.Earliest,
            EnableAutoCommit = false,
        };
        using var consumer = new ConsumerBuilder<string, string>(cfg).Build();   // <string,string> = key และ value ของ Kafka message ตัวอย่างนี้เป็น string ทั้งคู่
        consumer.Subscribe("user.events");

        try
        {
            while (!ct.IsCancellationRequested)
            {
                try
                {
                    var cr = consumer.Consume(ct);   // blocking sync
                    await Process(cr.Message.Value, ct);
                    consumer.Commit(cr);             // manual commit หลัง process สำเร็จ
                }
                catch (OperationCanceledException) { break; }
                catch (Exception ex)
                {
                    _logger.LogError(ex, "kafka consume failed");
                    // กรณีต้องการ retry: seek กลับ offset เดิม หรือใช้ DLQ topic
                }
            }
        }
        finally
        {
            consumer.Close();   // commit final offsets + leave group
        }
    }
}

7. Cron Schedule — Quartz.NET

BackgroundService + Task.Delay ทำงานตามรอบง่าย ๆ ได้

แต่ถ้าต้องการ feature ระดับ scheduler จริง ๆ ควรใช้ Quartz.NET:

  • cron expression — รูปแบบกำหนดเวลาทำงานเป็นรอบ เช่น "ทุกวัน 2:00", "ทุก 15 นาที"
  • misfire handling — จัดการเมื่อพลาดรอบที่ควรทำ (เช่นแอปดับอยู่ตอนถึงเวลา) ว่าจะข้ามไป หรือทำชดเชย
bash
dotnet add package Quartz.Extensions.Hosting
csharp
builder.Services.AddQuartz(q =>
{
    var key = new JobKey("CleanupJob");
    q.AddJob<CleanupJob>(o => o.WithIdentity(key));
    q.AddTrigger(o => o
        .ForJob(key)
        .WithIdentity("CleanupTrigger")
        .WithCronSchedule("0 0 2 * * ?"));   // ทุกวัน 2:00
});
builder.Services.AddQuartzHostedService();

public class CleanupJob : IJob
{
    public async Task Execute(IJobExecutionContext ctx) { /*...*/ }
}

Quartz feature:

  • Cron expression
  • Persistent store (DB) → schedule รอด restart
  • Distributed cluster (1 job 1 leader)
  • Misfire policy (ถ้า miss schedule, ทำหรือไม่)

8. Hangfire — Job เก็บใน DB + Dashboard

Quartz schedule ได้ดีแต่ job หายเมื่อแอป restart — Hangfire เก็บ job ลง DB ทำให้:

  • job รอดการ restart
  • retry อัตโนมัติเมื่อ fail
  • มี Dashboard (หน้าจอดูสถานะ job ผ่าน web)

เหมาะกับงานที่ต้องการความทนทานและมองเห็นได้

ศัพท์ Hangfire 3 แบบ:

  • enqueue — ใส่งานเข้าคิวทำทันที
  • schedule — ตั้งเวลาทำครั้งเดียวในอนาคต
  • recurring — ทำซ้ำเป็นรอบตาม cron
bash
dotnet add package Hangfire.AspNetCore
dotnet add package Hangfire.PostgreSql
csharp
builder.Services.AddHangfire(c => c.UsePostgreSqlStorage(connStr));
builder.Services.AddHangfireServer();

app.UseHangfireDashboard("/hangfire");

// enqueue
BackgroundJob.Enqueue<EmailService>(s => s.Send("a@b", "hi"));

// schedule
BackgroundJob.Schedule<EmailService>(s => s.Send("a@b", "hi"), TimeSpan.FromMinutes(5));

// recurring
RecurringJob.AddOrUpdate<CleanupService>("cleanup", s => s.Clean(), Cron.Daily);

ดี: Dashboard ดู job status, retry built-in, persistent


9. SignalR — Realtime push

เมื่อ background service ประมวลผลเสร็จแล้วอยาก push ผลให้ client เห็นทันที (เช่นแจ้งเตือน, progress bar) SignalR คือคำตอบ — มัน abstract WebSocket + fallback + group + reconnect ให้หมด เขียนแค่ hub แล้ว push เข้า group ที่ต้องการ:

bash
dotnet add package Microsoft.AspNetCore.SignalR
csharp
public class ChatHub : Hub
{
    public Task Send(string room, string msg) =>
        Clients.Group(room).SendAsync("message", msg);

    public Task Join(string room) => Groups.AddToGroupAsync(Context.ConnectionId, room);
}

builder.Services.AddSignalR();
app.MapHub<ChatHub>("/chat");

ผูก WebSocket + fallback + group + reconnect ให้ทั้งหมด


10. Worker Service (standalone, no web)

ถ้างานเป็น background ล้วน ๆ (consumer, cron) ไม่ต้องมี HTTP server เลย ใช้ template Worker Service ของ .NET ได้

แตกประเด็น:

  • dotnet new worker สร้างแอปที่มีแต่ host + hosted service
  • ไม่มี Kestrel (web server ของ ASP.NET Core) → process เบาลง ใช้ memory น้อยลง
  • เหมาะ deploy เป็น sidecar = container ผู้ช่วยที่รันคู่ไปกับ container หลักใน pod เดียวกัน (รูปแบบของ Kubernetes)
  • หรือ deploy เป็น consumer/cron container แยกจาก web API
bash
dotnet new worker -n MyWorker
csharp
// Program.cs
var builder = Host.CreateApplicationBuilder(args);
builder.Services.AddHostedService<Worker>();
var host = builder.Build();
host.Run();

ไม่มี Kestrel — แค่ background job — ใช้ deploy เป็น sidecar / consumer / cron container


11. Health Check ของ Background Service

ปัญหาที่อันตรายของ background service คือ อาการหยุดทำงานแบบเงียบ (silent stall) — process ยังอยู่ในระบบ ports ยังตอบ HTTP แต่ worker loop ค้างไม่ทำงานอยู่ข้างใน

ทำไม health check ปกติจับไม่เจอ:

  • HTTP health check ตอบ 200 ตามปกติ (เพราะ web server ยังรัน)
  • ตัวกลางเข้าใจว่าแอปยังสุขภาพดี
  • แต่ที่จริง worker ไม่ขยับมานานแล้ว

ทางแก้:

  1. ให้ worker อัปเดต "heartbeat" (สัญญาณว่ายังมีชีวิต) ทุกรอบที่ทำงาน
  2. health check อ่าน heartbeat แล้วถือว่า unhealthy ถ้าค้างนานเกินกำหนด
  3. k8s (Kubernetes — ระบบ orchestrate container) เห็น unhealthy แล้ว restart pod อัตโนมัติ
csharp
// ⚠️ ใช้ instance field (ไม่ใช่ static) เพื่อให้แต่ละ worker มี health state ของตัวเอง
// register เป็น Singleton เพื่อให้ worker และ health check ใช้ instance เดียวกัน
public class WorkerHealthCheck : IHealthCheck
{
    private volatile bool _isHealthy = true;
    public void MarkUnhealthy() => _isHealthy = false;
    public Task<HealthCheckResult> CheckHealthAsync(HealthCheckContext c, CancellationToken ct) =>
        Task.FromResult(_isHealthy ? HealthCheckResult.Healthy() : HealthCheckResult.Unhealthy());
}

// register Singleton ให้ worker และ health check ใช้ instance เดียวกัน
services.AddSingleton<WorkerHealthCheck>();
services.AddHealthChecks().AddCheck<WorkerHealthCheck>("worker");

Worker จะ mark unhealthy เมื่อ stuck — k8s restart


12. Distributed Lock (multiple replicas)

กับดักสำคัญเมื่อ scale หลาย replica — ถ้าทุก replica รัน background job เดียวกัน job จะถูกทำซ้ำ (เช่นส่ง email ซ้ำ 3 รอบ)

วิธีแก้ใช้ distributed lock — กลอนล็อกที่ใช้ร่วมกันข้ามหลายเครื่อง

แตกประเด็น:

  • replica ทุกตัวพยายามจอง lock ก่อนทำ job
  • มีตัวเดียวที่ได้ lock → ตัวนั้นเป็นคนทำ
  • ตัวอื่นที่จองไม่ได้ → ข้ามรอบนี้ไป
  • ตัวอย่าง: Redis RedLock, etcd lease, ZooKeeper

ถ้ามี 3 replica + background job → ทุก replica จะรัน job → ต้อง coordinate

Redis-based lock (StackExchange.Redis + RedLockNet):

csharp
using RedLockNet.SERedis;

var multiplexer = ConnectionMultiplexer.Connect("redis:6379");
var factory = RedLockFactory.Create(new[] { multiplexer });

await using var redLock = await factory.CreateLockAsync("job-cleanup", TimeSpan.FromMinutes(10));
if (redLock.IsAcquired)
{
    await DoJob();
}

⚠️ ข้อจำกัดของ RedLock — RedLock ไม่รับประกัน correctness (ความถูกต้องแน่นอน) ใต้สถานการณ์ที่หนักหน่วง เช่น network partition (เครือข่ายขาดการเชื่อมต่อบางส่วน แต่ละฝั่งเห็นสถานะกันไม่ครบ) หรือ clock drift (นาฬิกาของแต่ละเครื่องคลาดเคลื่อนไม่ตรงกัน) รุนแรง ๆ — นี่คือประเด็นที่ยังถกเถียงกันในวงการ ไม่ใช่ข้อสรุปที่ทุกคนเห็นตรงกัน (ผู้เขียนหนังสือ "Designing Data-Intensive Applications" เคยเขียนวิจารณ์ RedLock ไว้ และผู้เขียน Redis เองก็เคยตอบโต้กลับมา) สำหรับงานสำคัญที่ห้ามทำซ้ำเด็ดขาด (เช่น financial transaction) ทางเลือกที่ปลอดภัยกว่าคือใช้ leader election (การเลือกตัวแทนเดียวให้เป็นคนทำงานแทนทั้งกลุ่ม) ผ่าน etcd, Consul, หรือ ZooKeeper แทน ส่วน RedLock ยังเหมาะกับงาน coordination ทั่วไปที่ต่อให้ทำซ้ำบ้างก็ไม่เสียหายมาก


13. Outbox Pattern (กับ event-driven)

outbox pattern แก้ปัญหา "บันทึก DB สำเร็จ แต่ publish event ไป Kafka ล้มเหลว (หรือกลับกัน)" — สถานะระบบไม่ตรงกัน

วิธีทำ:

  1. แทนที่จะ publish event ตรง ๆ → เขียน event ลงตาราง OutboxMessages
  2. เขียนใน transaction เดียวกับ business data (ข้อมูลธุรกิจ เช่นออเดอร์, ผู้ใช้)
  3. ทั้งคู่จึง atomic — สำเร็จพร้อมกันหรือ rollback พร้อมกัน ไม่มีกรณีฝั่งใดฝั่งหนึ่งสำเร็จเดี่ยว ๆ
  4. มี background service คอย poll ตาราง OutboxMessages แล้ว publish ไป Kafka ทีหลัง พร้อม mark ว่าส่งแล้ว
csharp
// ภายใน transaction ของ business operation
// (event-driven = สถาปัตยกรรมที่ส่วนต่าง ๆ สื่อสารกันผ่าน event)
await using var tx = await db.Database.BeginTransactionAsync();
db.Orders.Add(order);
db.OutboxMessages.Add(new OutboxMessage { Topic = "order.created", Payload = JsonSerializer.Serialize(order) });
await db.SaveChangesAsync();
await tx.CommitAsync();

// background service poll outbox table → publish ไป Kafka → mark sent

รับประกัน "publish event เกิดขึ้น atomic กับ DB write"

ข้อควรรู้เพิ่ม:

  • ฝั่ง consumer ควรออกแบบให้ idempotent (ทำซ้ำ event ตัวเดิมแล้วได้ผลเดิม) เพราะ outbox อาจ publish ซ้ำตอน retry — เก็บ message id แล้ว skip ถ้าเคยทำแล้ว
  • ทางเลือกอีกแบบคือ CDC-based outbox (Change Data Capture) — ใช้ tool เช่น Debezium อ่าน WAL/binlog ของ DB ส่งเข้า Kafka ให้เอง ไม่ต้องเขียน background poller

14. Graceful Shutdown ของ BackgroundService

ตอนแอปกำลังปิด background service ที่กำลังประมวลผล job อยู่ต้องหยุดอย่างสะอาด — ไม่ทิ้งงานค้างกลางคัน

กุญแจคือต้อง honor (เคารพ/ทำตาม) CancellationToken ที่ framework ส่งมา:

  • token = สัญญาณ "ขอให้ยกเลิกงาน" ที่ส่งผ่าน async chain
  • ส่งต่อให้ทุก async call ที่เรียก (DB, HTTP, queue ฯลฯ)
  • ทุก call จะรับรู้สัญญาณยกเลิกแล้วหยุดอย่างเรียบร้อย
  • ทำให้ job ปัจจุบันมีโอกาส commit, rollback, หรือทำ cleanup ก่อน process จะถูกฆ่า
csharp
public class JobWorker : BackgroundService
{
    private readonly ILogger<JobWorker> _logger;
    public JobWorker(ILogger<JobWorker> logger) => _logger = logger;

    protected override async Task ExecuteAsync(CancellationToken ct)
    {
        try
        {
            while (!ct.IsCancellationRequested)
            {
                var job = await DequeueAsync(ct);
                await ProcessAsync(job, ct);   // ต้อง honor ct ทุก async call
            }
        }
        catch (OperationCanceledException) when (ct.IsCancellationRequested)
        {
            _logger.LogInformation("JobWorker stopping gracefully");
        }
    }

    public override async Task StopAsync(CancellationToken ct)
    {
        _logger.LogInformation("JobWorker stop signal received");
        // จุดนี้ ExecuteAsync จะเห็น ct fire แล้ว
        // base.StopAsync รอ ExecuteAsync จบ ภายในเวลา HostOptions.ShutdownTimeout
        await base.StopAsync(ct);
        // ปิด resource ที่เหลือตรงนี้ (connection, file handle ฯลฯ)
    }
}

💡 HostOptions.ShutdownTimeout (default 30 วินาที) = กรอบเวลาที่ host จะรอ StopAsync ของทุก service จนจบ ถ้าเกินเวลา host จะเดินหน้าปิด process ต่อไป — งานที่ยังค้างจะถูก kill จึงต้องออกแบบ job ให้ทำเสร็จเร็วหรือ checkpoint ได้


15. Common Pitfalls

  1. Block (ค้าง) ใน StartAsync → app startup ช้า / fail
  2. ลืม try/catch ใน ExecuteAsync — exception หลุด → service ตาย → host crash (ทั้งแอปดับ)
  3. หลาย replica + job ไม่ lock → รันซ้ำ
  4. Sync over async (เรียกโค้ด async แบบบล็อกรอผล) ใน hosted service — บล็อก thread
  5. ใช้ scoped service ตรง ๆ ใน BackgroundService → captive dependency — ต้อง inject IServiceScopeFactory แล้ว CreateAsyncScope()
  6. ไม่ honor CancellationToken → ปิดตัวไม่ได้ → k8s force kill (สั่งฆ่าแบบบังคับ)
  7. Channel unbounded + producer เร็ว → OOM (Out Of Memory — หน่วยความจำเต็มจนแอป crash)
  8. loop ทำงาน synchronous ยาวไม่หยุดพัก (ไม่มีจุด yield คืน thread ให้งานอื่น) → ผลคือ CPU starve คืองานอื่นอดใช้ CPU

🛠️ Checkpoint 8 — ลงมือทำ

  1. BackgroundService ที่ลบ log เก่า > 30 วัน ทุก 1 ชม.
  2. queue worker ด้วย Channel + producer ใน API endpoint
  3. Quartz cron job ที่รันทุก midnight + persist ใน PG
  4. Hangfire dashboard + ลอง enqueue + retry
  5. SignalR hub สำหรับ chat broadcast (group)
  6. distributed lock — 3 replica + cleanup job ทำงาน 1 ตัว
  7. test: kill pod ขณะ ExecuteAsync → graceful exit (ct fire)

สรุปบทที่ 8 + ปิดเล่ม

ครบเล่ม .NET แล้ว:

  1. ASP.NET Core pipeline + DI
  2. Minimal + Controller API
  3. EF Core + Postgres + migration
  4. Auth (JWT + refresh + RBAC + policy + resource-based)
  5. Test (unit + integration + e2e + Testcontainers)
  6. Production (Serilog, health, OpenTelemetry, Docker, AOT)
  7. Middleware + filter + model binding ลึก
  8. DI + Options + HttpClient + Polly
  9. Background service + Channels + Queue + Quartz + Hangfire

ขั้นถัดไป (หัวข้อที่ควรเรียนต่อนอกเล่มนี้):

  • microservices — pattern เชิงสถาปัตยกรรม (saga, CQRS, event-driven, service mesh)
  • observability — logging / tracing / metrics / SLO ลึก (OpenTelemetry, Prometheus, Grafana, Tempo)
  • system design — scale ใหญ่ (sharding, replication, caching strategy, queue capacity planning)
  • security — defensive engineering (threat modeling, secret management, supply chain security)

ขอให้สนุกกับ .NET!

กลับสารบัญหลัก