종합 실습 - 배송 추적 처리 파이프라인
이 장에서 배우는 것
앞 장까지 제네릭, 델리게이트, 반복자, 메모리, 비동기, 동시성, 의존성 주입, 테스트 설계, 리플렉션, 패턴 매칭, 성능 측정을 하나씩 다뤘다. 이 장에서는 그 가운데 비동기 처리와 동시성 도구만 골라 택배 배송 추적 이벤트를 처리하는 파이프라인 한 벌을 만든다. 새 문법은 나오지 않는다. 이미 아는 부품을 어떤 순서로 조립하고, 어디에서 멈추고, 실패를 어디까지 퍼지게 둘지 정하는 일이 이 장의 주제다.
파이프라인은 접수, 파싱, 검증, 집계의 네 단계로 구성한다. 단계 사이는 Channel 로 잇고, 도중에 취소해도 프로그램이 멈추지 않고 끝나게 하며, 잘못된 입력 한 줄이 전체를 무너뜨리지 않게 한다. 마지막에는 결과를 집계해 보고서로 출력하고, 설계를 되돌아본다.
- bounded 채널로 단계를 잇고 단계마다 worker 수를 다르게 두는 구조를 만든다.
- 취소 요청이 모든 단계로 퍼지고, 각 단계가 다음 채널을 닫으며 끝나는 흐름을 구현한다.
- 항목 하나의 오류를 격리해 실패 목록에 모으고, 파이프라인은 계속 돌리는 방법을 익힌다.
- worker 의 완료 순서가 바뀌어도 같은 보고서가 나오도록 집계 규칙을 정한다.
- 완성한 설계를 회고하며 다음에 바꿀 지점을 짚는다.
문제 상황
택배 물류 센터에는 배송 추적 이벤트가 쉬지 않고 들어온다. 접수, 분류, 출고, 배송 완료 같은 상태 변화가 한 줄 문자열로 도착한다. 처음에는 foreach 로 한 줄씩 읽어 파싱하고, 허브를 조회하고, 딕셔너리에 반영하는 코드로 충분했다. 물량이 늘면서 세 가지 문제가 한꺼번에 나타났다.
첫째, 허브 조회는 네트워크를 타는 작업이라 한 건씩 기다리면 전체가 느려진다. 그렇다고 이벤트 수만큼 Task 를 띄우면 입력이 몰릴 때 메모리와 연결 수가 한없이 늘어난다. 둘째, 형식이 깨진 줄이나 존재하지 않는 허브 코드가 섞여 들어온다. 예외 하나가 루프 전체를 끊으면 뒤에 있는 정상 이벤트까지 처리되지 않는다. 셋째, 운영자가 작업을 중단시키면 진행 중이던 일이 어중간한 상태로 남는다. 어디까지 처리했는지, 남은 것은 무엇인지 알 수 없다.
이 세 가지는 서로 얽혀 있다. 동시성을 늘리면 순서가 흔들려 집계 결과가 실행마다 달라질 수 있고, 취소를 넣으면 채널을 누가 닫는지가 모호해져 종료되지 않는 프로그램이 되기 쉽다. 그래서 단계 구조, 종료 규칙, 오류 규칙, 집계 규칙을 코드를 쓰기 전에 먼저 정한다.
단계 구조와 채널
파이프라인은 단계를 채널로 잇는 구조다. 앞 단계는 결과를 채널에 쓰고, 뒤 단계는 채널에서 읽는다. 단계 하나는 같은 일을 하는 worker 여러 개로 구성할 수 있다. 느린 단계에는 worker 를 늘리고 빠른 단계는 하나만 둔다.
채널은 반드시 용량을 제한한 bounded 형태로 만든다. 용량이 가득 차면 WriteAsync 가 자리가 날 때까지 기다리므로, 느린 단계가 앞 단계의 속도를 자연스럽게 늦춘다. 이것을 배압(backpressure)이라 부른다. 무제한 채널은 코드가 짧지만 소비가 생산을 따라가지 못하면 메모리가 계속 늘어난다.
| 단계 | worker | 입력 → 출력 | 실패 처리 |
|---|---|---|---|
| 접수 | 1 | 문자열 목록 → RawLine | 없음 |
| 파싱 | 2 | RawLine → TrackingEvent | 형식 오류를 실패 목록에 기록 |
| 검증 | 3 | TrackingEvent → CheckedEvent | 허브 조회 오류를 실패 목록에 기록 |
| 집계 | 1 | CheckedEvent → 보고서 | 없음 |
집계 단계 worker 를 하나로 둔 데는 이유가 있다. 집계용 딕셔너리를 한 스레드만 만지게 하면 lock 이 필요 없다. 공유 상태를 보호하는 대신 공유 상태에 접근하는 코드를 한 곳으로 모은 것이다.
종료·취소·오류를 다루는 규칙
파이프라인이 끝나는 길은 셋이다. 입력이 다 소진되어 정상 종료하는 길, 취소 요청으로 중단하는 길, 예상하지 못한 예외로 실패하는 길이다. 세 길 모두에서 모든 단계가 끝나야 하고, 아무도 무한히 기다리면 안 된다. 이 장의 코드는 다음 규칙으로 이를 보장한다.
| 경로 | 시작 신호 | 단계의 반응 | 다음 채널 |
|---|---|---|---|
| 정상 종료 | 접수 단계가 입력을 소진 | 읽기 루프가 자연스럽게 끝남 | 정상 완료 |
| 취소 | cts.Cancel() | 대기 중인 작업이 OperationCanceledException 으로 깨어남 | 예외와 함께 완료 |
| 예상 밖 오류 | worker 에서 예외 발생 | 토큰을 취소해 나머지 단계도 멈춤 | 예외와 함께 완료 |
핵심은 채널을 닫는 책임이 worker 개별이 아니라 단계 전체에 있다는 점이다. worker 가 셋인 단계에서 첫 번째 worker 가 일을 끝내고 Complete 를 부르면, 아직 쓰고 있는 나머지 worker 의 쓰기가 실패한다. 그래서 단계의 모든 worker 가 끝난 뒤 한 번만 닫는다. 이 일을 StageAsync 하나가 맡는다.
오류는 두 종류로 나눈다. 입력 한 건이 나쁜 경우는 예상 가능한 오류다. 이 경우 그 건만 실패 목록에 넣고 다음 건으로 넘어간다. 반면 코드의 버그처럼 예상하지 못한 예외는 어느 항목에서나 다시 날 수 있으므로 파이프라인 전체를 취소하고 밖으로 다시 던진다. 항목 단위 catch 에서 OperationCanceledException 은 반드시 제외해야 한다. 이 예외까지 삼키면 취소가 무시된다.
결정적인 집계
worker 가 여럿이면 이벤트가 집계 단계에 도착하는 순서는 실행마다 다르다. 그래서 "가장 나중에 도착한 상태가 최신 상태"라는 규칙은 쓸 수 없다. 이 장에서는 상태 열거형이 진행 순서(접수, 분류, 출고, 배송 완료)대로 정의되어 있다는 점을 이용해 더 큰 상태 값을 최신으로 본다. 도착 순서가 바뀌어도 결과는 같다. 출력할 때는 키를 정렬하고, 실패 목록은 입력 번호순으로 정렬한다.
이렇게 하면 실행 결과를 문서에 그대로 실을 수 있다. 시간이나 스레드 순서에 따라 달라지는 값은 아예 출력하지 않는다.
완성 코드
아래 프로그램은 Program.cs 한 파일이다. 프로젝트는 dotnet new console 로 만든 .NET 10 콘솔 앱이면 된다. 입력 15줄 가운데 5줄은 일부러 잘못된 값이다. 전체 처리를 한 번, 6번째 줄까지 접수한 뒤 취소하는 처리를 한 번 실행한다.
using System.Collections.Concurrent;
using System.Threading.Channels;
string[] lines =
[
"P001,Received,SEOUL",
"P002,Received,BUSAN",
"P001,Sorted,SEOUL",
"P003,Received,DAEJEON",
"P002,Sorted,BUSAN",
"P001,Shipped,SEOUL",
"P004,Lost,SEOUL",
"P003,Sorted,DAEJEON",
"P005,Received,MARS",
"P001,Delivered,GWANGJU",
"broken-line",
"P002,Shipped,BUSAN",
"P003,Shipped,ATLANTIS",
"P006,Received,INCHEON",
",Sorted,SEOUL",
];
var full = await new TrackingPipeline().RunAsync(lines, int.MaxValue);
Console.WriteLine("== 전체 처리 ==");
Console.WriteLine($"접수 {full.Accepted}건, 집계 {full.Processed}건, 실패 {full.Failures.Count}건, 취소 {YesNo(full.Canceled)}");
Console.WriteLine("[실패 목록]");
foreach (var f in full.Failures)
Console.WriteLine($" #{f.Seq} {f.Stage}: {f.Reason}");
Console.WriteLine("[권역별 이벤트]");
foreach (var (region, count) in full.RegionCounts.OrderBy(p => p.Key, StringComparer.Ordinal))
Console.WriteLine($" {region}: {count}");
Console.WriteLine("[택배별 최종 상태]");
foreach (var (id, status) in full.Latest.OrderBy(p => p.Key, StringComparer.Ordinal))
Console.WriteLine($" {id}: {status}");
var cut = await new TrackingPipeline().RunAsync(lines, 6);
Console.WriteLine("== 중간 취소 ==");
Console.WriteLine($"접수 {cut.Accepted}건, 취소 {YesNo(cut.Canceled)}");
var total = cut.Processed + cut.Failures.Count;
Console.WriteLine($"집계+실패 합계가 접수 건수 이하: {YesNo(total <= cut.Accepted)}");
static string YesNo(bool value) => value ? "예" : "아니오";
enum Status { Received, Sorted, Shipped, Delivered }
sealed record RawLine(int Seq, string Text);
sealed record TrackingEvent(int Seq, string ParcelId, Status Status, string Hub);
sealed record CheckedEvent(TrackingEvent Event, string Region);
sealed record Failure(int Seq, string Stage, string Reason);
sealed record PipelineResult(
bool Canceled,
int Accepted,
int Processed,
IReadOnlyList<Failure> Failures,
IReadOnlyDictionary<string, int> RegionCounts,
IReadOnlyDictionary<string, Status> Latest);
sealed class TrackingPipeline
{
static readonly Dictionary<string, string> Regions = new()
{
["SEOUL"] = "수도권",
["INCHEON"] = "수도권",
["BUSAN"] = "영남",
["DAEJEON"] = "충청",
["GWANGJU"] = "호남",
};
readonly ConcurrentQueue<Failure> _failures = new();
readonly Dictionary<string, int> _regionCounts = new();
readonly Dictionary<string, Status> _latest = new();
int _accepted;
int _processed;
public async Task<PipelineResult> RunAsync(IReadOnlyList<string> input, int stopAfter)
{
using var cts = new CancellationTokenSource();
var ct = cts.Token;
var rawCh = Channel.CreateBounded<RawLine>(4);
var parsedCh = Channel.CreateBounded<TrackingEvent>(4);
var checkedCh = Channel.CreateBounded<CheckedEvent>(4);
var stages = new[]
{
StageAsync(cts, 1, e => rawCh.Writer.TryComplete(e),
() => ProduceAsync(input, rawCh.Writer, stopAfter, cts)),
StageAsync(cts, 2, e => parsedCh.Writer.TryComplete(e),
() => ParseWorkerAsync(rawCh.Reader, parsedCh.Writer, ct)),
StageAsync(cts, 3, e => checkedCh.Writer.TryComplete(e),
() => ValidateWorkerAsync(parsedCh.Reader, checkedCh.Writer, ct)),
StageAsync(cts, 1, _ => { },
() => AggregateAsync(checkedCh.Reader, ct)),
};
try
{
await Task.WhenAll(stages);
}
catch (OperationCanceledException)
{
}
return new PipelineResult(
cts.IsCancellationRequested,
_accepted,
_processed,
_failures.OrderBy(f => f.Seq).ToList(),
_regionCounts,
_latest);
}
static async Task StageAsync(
CancellationTokenSource cts, int workers, Action<Exception?> complete, Func<Task> body)
{
Exception? error = null;
try
{
await Task.WhenAll(Enumerable.Range(0, workers).Select(_ => body()));
}
catch (OperationCanceledException ex)
{
error = ex;
throw;
}
catch (Exception ex)
{
error = ex;
cts.Cancel();
throw;
}
finally
{
complete(error);
}
}
async Task ProduceAsync(
IReadOnlyList<string> input, ChannelWriter<RawLine> writer, int stopAfter, CancellationTokenSource cts)
{
for (var i = 0; i < input.Count; i++)
{
if (i == stopAfter)
{
cts.Cancel();
return;
}
await writer.WriteAsync(new RawLine(i + 1, input[i]), cts.Token);
_accepted++;
}
}
async Task ParseWorkerAsync(
ChannelReader<RawLine> input, ChannelWriter<TrackingEvent> output, CancellationToken ct)
{
await foreach (var raw in input.ReadAllAsync(ct))
{
var (ev, error) = Parse(raw);
if (ev is not null)
await output.WriteAsync(ev, ct);
else
_failures.Enqueue(new Failure(raw.Seq, "parse", error ?? "알 수 없는 오류"));
}
}
static (TrackingEvent? Event, string? Error) Parse(RawLine raw)
{
var parts = raw.Text.Split(',');
if (parts.Length is not 3)
return (null, $"필드 수 {parts.Length}개 (3개 필요)");
var id = parts[0].Trim();
if (id.Length == 0)
return (null, "택배 번호가 비어 있음");
var statusText = parts[1].Trim();
if (Enum.TryParse<Status>(statusText, out var status) && Enum.IsDefined(status))
return (new TrackingEvent(raw.Seq, id, status, parts[2].Trim()), null);
return (null, $"알 수 없는 상태: {statusText}");
}
async Task ValidateWorkerAsync(
ChannelReader<TrackingEvent> input, ChannelWriter<CheckedEvent> output, CancellationToken ct)
{
await foreach (var ev in input.ReadAllAsync(ct))
{
string region;
try
{
region = await ResolveRegionAsync(ev.Hub, ct);
}
catch (Exception ex) when (ex is not OperationCanceledException)
{
_failures.Enqueue(new Failure(ev.Seq, "validate", ex.Message));
continue;
}
await output.WriteAsync(new CheckedEvent(ev, region), ct);
}
}
static async Task<string> ResolveRegionAsync(string hub, CancellationToken ct)
{
await Task.Delay(1, ct);
return Regions.TryGetValue(hub, out var region)
? region
: throw new InvalidOperationException($"알 수 없는 허브: {hub}");
}
async Task AggregateAsync(ChannelReader<CheckedEvent> input, CancellationToken ct)
{
await foreach (var item in input.ReadAllAsync(ct))
{
var ev = item.Event;
_processed++;
_regionCounts[item.Region] = _regionCounts.GetValueOrDefault(item.Region) + 1;
if (_latest.TryGetValue(ev.ParcelId, out var old) is false || ev.Status > old)
_latest[ev.ParcelId] = ev.Status;
}
}
}
줄별 해설
입력과 출력부
최상위 문장은 입력 배열을 만들고 파이프라인을 두 번 돌린다. new TrackingPipeline() 을 실행마다 새로 만드는 이유는 실패 큐와 집계 딕셔너리가 인스턴스 필드이기 때문이다. 한 인스턴스는 한 번만 실행한다는 전제다. RunAsync 의 두 번째 인자는 "몇 건을 접수한 뒤 취소할지"다. int.MaxValue 는 취소하지 않는다는 뜻이다. 출력할 때 OrderBy(..., StringComparer.Ordinal) 로 키를 정렬하는 것은 딕셔너리 순회 순서에 기대지 않기 위해서다.
RunAsync 와 StageAsync
RunAsync 는 채널 세 개를 용량 4로 만들고, 단계 네 개를 배열로 시작한 뒤 Task.WhenAll 로 모두 기다린다. 각 단계의 두 번째 숫자가 worker 수다. 세 번째 인자 e => rawCh.Writer.TryComplete(e) 는 "이 단계가 끝나면 다음 채널을 이 오류와 함께 닫아라"는 함수다. 마지막 집계 단계는 닫을 채널이 없어 _ => { } 를 넘긴다.
StageAsync 는 worker 를 workers 개 만들어 모두 기다린다. 취소 예외는 기록만 하고 다시 던지며, 그 밖의 예외는 cts.Cancel() 로 다른 단계까지 멈춘 뒤 다시 던진다. 어느 경로든 finally 에서 complete(error) 가 실행되므로 다음 단계는 반드시 깨어난다. 오류와 함께 닫힌 채널은 읽는 쪽에서 그 예외를 만나 끝난다.
WhenAll 을 감싼 catch (OperationCanceledException) 는 비어 있다. 취소는 오류가 아니라 정상적인 종료 경로이기 때문이다. 결과의 Canceled 는 예외를 잡았는지가 아니라 cts.IsCancellationRequested 로 판단한다. 예외 전달 경로의 미세한 차이에 결과가 좌우되지 않게 하려는 것이다.
접수 단계
ProduceAsync 는 번호를 붙인 RawLine 을 채널에 쓴다. 채널이 가득 차 있으면 WriteAsync 가 기다린다. i == stopAfter 이면 스스로 취소를 요청하고 돌아온다. 실제 서비스에서는 이 자리에 외부의 취소 요청이 들어온다. 접수 건수는 쓰기가 끝난 뒤에 세므로, _accepted 는 채널에 실제로 들어간 건수다. 접수 단계 worker 는 하나뿐이라 증가 연산에 동기화가 필요 없다.
파싱과 검증 단계
Parse 는 예외를 던지지 않고 값 또는 사유를 튜플로 돌려준다. 형식이 깨진 입력은 예외 상황이 아니라 흔한 상황이기 때문이다. 필드 수, 빈 택배 번호, 정의되지 않은 상태의 세 검사가 순서대로 실행된다. Enum.IsDefined 검사를 두는 것은 숫자 문자열이 정의 밖의 값으로 변환되는 경우를 막기 위해서다.
검증 단계는 허브를 비동기로 조회한다. 예제에서는 Task.Delay(1, ct) 로 네트워크 지연을 흉내 낸다. 조회에서 예외가 나면 when (ex is not OperationCanceledException) 조건으로 취소를 제외하고 잡아 실패 목록에 넣은 뒤 continue 로 다음 항목으로 간다. try 범위를 조회 한 줄로 좁힌 것도 의도다. 채널 쓰기에서 나는 예외가 "항목 실패"로 잘못 분류되지 않는다.
집계 단계
집계는 단일 worker 가 수행한다. 권역별 건수를 세고, 택배별로 더 큰 상태 값을 최신으로 갱신한다. TryGetValue(...) is false || ev.Status > old 는 "처음 보는 택배이거나 더 진행된 상태이면 갱신한다"는 조건이다. 이벤트가 어떤 순서로 도착해도 결과가 같다. 예를 들어 P003 의 출고 이벤트는 허브 코드가 잘못되어 검증에서 탈락했으므로, P003 의 최종 상태는 분류로 남는다.
실행 결과
프로젝트 폴더에서 다음 명령을 실행한다.
$ dotnet run
== 전체 처리 ==
접수 15건, 집계 10건, 실패 5건, 취소 아니오
[실패 목록]
#7 parse: 알 수 없는 상태: Lost
#9 validate: 알 수 없는 허브: MARS
#11 parse: 필드 수 1개 (3개 필요)
#13 validate: 알 수 없는 허브: ATLANTIS
#15 parse: 택배 번호가 비어 있음
[권역별 이벤트]
수도권: 4
영남: 3
충청: 2
호남: 1
[택배별 최종 상태]
P001: Delivered
P002: Shipped
P003: Sorted
P006: Received
== 중간 취소 ==
접수 6건, 취소 예
집계+실패 합계가 접수 건수 이하: 예
취소 실행에서 집계 건수와 실패 건수는 각각 0부터 6 사이에서 실행마다 달라질 수 있다. 그래서 값 자체는 출력하지 않고 항상 참이어야 하는 관계만 확인한다. 이 실행의 접수 건수 6은 접수 단계가 6번째 줄까지 쓴 뒤 스스로 취소하므로 항상 같다.
실무에서 자주 틀리는 것
무제한 채널로 배압을 없앤다
틀린 코드는 다음과 같다.
var rawCh = Channel.CreateUnbounded<RawLine>();
접수가 검증보다 빠르면 채널에 이벤트가 쌓이며 메모리가 늘어난다. 평소에는 티가 나지 않다가 입력이 몰리는 시간대에 문제가 된다. 용량을 정해 두면 생산자가 알아서 속도를 늦춘다.
var rawCh = Channel.CreateBounded<RawLine>(4);
worker 마다 채널을 닫는다
틀린 코드는 worker 가 끝날 때 각자 Complete 를 부른다.
async Task WorkerAsync(ChannelReader<A> input, ChannelWriter<B> output)
{
await foreach (var a in input.ReadAllAsync())
await output.WriteAsync(Convert(a));
output.Complete();
}
가장 먼저 끝난 worker 가 채널을 닫으면 아직 쓰는 중인 worker 의 WriteAsync 가 ChannelClosedException 을 던진다. 단계의 모든 worker 를 Task.WhenAll 로 기다린 뒤 한 번만 닫는다. 이 장의 StageAsync 가 그 구조다.
try { await Task.WhenAll(workers); }
finally { output.TryComplete(); }
항목 오류를 잡다가 취소까지 삼킨다
틀린 코드는 다음과 같다.
try
{
region = await ResolveRegionAsync(ev.Hub, ct);
}
catch (Exception ex)
{
_failures.Enqueue(new Failure(ev.Seq, "validate", ex.Message));
continue;
}
취소로 발생한 OperationCanceledException 도 항목 실패로 기록되고 루프가 계속된다. 취소한 뒤에도 실패 건수가 늘고, 결국 읽기 쪽에서 다시 예외가 나서야 끝난다. 취소 예외를 조건으로 제외한다.
catch (Exception ex) when (ex is not OperationCanceledException)
{
_failures.Enqueue(new Failure(ev.Seq, "validate", ex.Message));
continue;
}
도착 순서를 결과에 반영한다
틀린 코드는 마지막으로 도착한 이벤트를 최신 상태로 본다.
_latest[ev.ParcelId] = ev.Status;
검증 단계 worker 가 셋이면 같은 택배의 분류 이벤트가 접수 이벤트보다 늦게 도착하는 일이 생긴다. 실행할 때마다 최종 상태가 달라질 수 있고, 테스트가 간헐적으로 실패한다. 진행 순서로 비교하는 규칙을 둔다.
if (_latest.TryGetValue(ev.ParcelId, out var old) is false || ev.Status > old)
_latest[ev.ParcelId] = ev.Status;
설계 회고
완성한 파이프라인을 되돌아보면 좋은 점과 남은 과제가 함께 보인다. 좋은 점은 규칙이 몇 개로 정리된다는 것이다. 채널은 bounded 로 만들고, 단계가 채널을 닫고, 항목 오류는 격리하고, 집계는 도착 순서에 기대지 않는다. 이 네 규칙을 지키면 worker 수나 채널 용량을 바꿔도 결과가 같다. 앞 장에서 다룬 성능 측정 도구로 병목 단계를 찾은 다음, 그 단계의 worker 수만 조정하면 된다.
남은 과제도 있다. 지금 파이프라인은 한 번 실행하고 끝나는 배치 형태이고, 취소된 뒤에 어디까지 처리했는지는 건수로만 알 수 있다. 재시작하려면 처리한 마지막 번호를 저장해야 한다. 실패 목록은 메모리에만 있어서 프로세스가 죽으면 사라진다. 또 검증 단계는 실제 서비스 호출 대신 지연을 흉내 낸 것이라, 실제 호출을 넣으면 시간 제한과 재시도 정책이 필요하다. 이런 요소는 단계 함수를 인터페이스 뒤로 옮기면 테스트에서 가짜 객체로 바꿔 끼울 수 있다. 앞에서 다룬 의존성 주입과 시간 추상화가 여기에 이어진다.
| 결정 | 이유 | 바꾸면 생기는 일 |
|---|---|---|
| bounded 채널(용량 4) | 배압으로 메모리 사용량을 제한 | unbounded 이면 입력이 몰릴 때 메모리 증가 |
단계 단위 Complete | 쓰는 중인 worker 의 쓰기 실패 방지 | worker 마다 닫으면 간헐적 예외 |
| 항목 단위 오류 격리 | 나쁜 입력 한 건이 전체를 멈추지 않게 함 | 격리하지 않으면 정상 이벤트가 유실 |
| 단일 집계 worker | lock 없이 딕셔너리 보호 | 여러 개로 늘리면 동기화 필요 |
| 진행 순서로 최신 상태 판정 | 도착 순서와 무관한 결정적 결과 | 도착 순서에 의존하면 실행마다 결과가 다름 |
한눈에 보기
| 도구 | 쓴 곳 | 역할 | 주의할 점 |
|---|---|---|---|
Channel.CreateBounded | 단계 사이 | 배압이 있는 전달 통로 | 용량은 측정으로 정한다 |
ReadAllAsync(ct) | 각 worker | 채널이 닫힐 때까지 읽기 | 취소되면 예외로 끝난다 |
TryComplete(error) | StageAsync | 다음 단계를 깨워 종료시킴 | 단계당 한 번, finally 에서 |
CancellationTokenSource | 접수, 오류 시 | 모든 단계에 취소 전달 | 항목 catch 에서 취소 예외를 제외 |
ConcurrentQueue | 실패 목록 | 여러 worker 가 안전하게 기록 | 출력 전에 번호순 정렬 |
공식 문서에서는 System.Threading.Channels 안내에서 채널의 옵션과 종료 동작을 확인할 수 있다.
연습 문제
- 검증 단계의 worker 수를 3에서 1로, 채널 용량을 4에서 1로 바꿔 실행한다. 전체 처리 출력이 달라지는지 예상하고 이유를 설명하시오.
- 집계 단계에 상태별 최종 건수(Received, Sorted, Shipped, Delivered 순)를 출력하는 기능을 추가하시오. 이 프로그램의 입력에서 예상되는 출력도 적으시오.
- 이벤트 3번째 줄에서
ResolveRegionAsync가OperationCanceledException이 아닌 예외를 던지는 대신, 코드 버그로NullReferenceException이ParseWorkerAsync안에서 발생했다고 하자. 이 장의 구조에서는 어떤 일이 일어나는지 순서대로 설명하시오. - 사용자가 Ctrl+C 를 눌렀을 때도 파이프라인이 정상적으로 취소되도록
RunAsync를 고치시오.
정답과 해설
- 출력은 달라지지 않는다. worker 수와 채널 용량은 처리 속도와 완료 순서에만 영향을 주고, 집계는 도착 순서에 의존하지 않으며 실패 목록은 번호순으로 정렬해 출력하기 때문이다. 다만 용량이 1이면 앞 단계가 더 자주 기다리므로 총 시간은 늘어난다. 결과가 같은 것은 이 장에서 정한 집계 규칙 덕분이다.
- 집계 단계에서 마지막에
_latest.Values를 상태별로 센다. 예를 들면 다음과 같다.
이 입력에서는 P001 이 Delivered, P002 가 Shipped, P003 이 Sorted, P006 이 Received 이므로 네 상태가 모두 1건이다. 출력은foreach (var status in Enum.GetValues<Status>()) { var n = full.Latest.Values.Count(s => s == status); Console.WriteLine($" {status}: {n}"); }Received: 1,Sorted: 1,Shipped: 1,Delivered: 1순서다.Enum.GetValues는 정의된 순서로 돌려주므로 출력 순서가 고정된다. - 파싱 worker 의 예외는 항목
catch가 없으므로StageAsync의 일반catch (Exception)에 도달한다. 순서는 이렇다. 먼저cts.Cancel()이 호출되어 다른 단계의 대기 중인 읽기와 쓰기가 취소 예외로 깨어난다. 다음으로finally에서 파싱 단계가 다음 채널을 예외와 함께 닫는다. 그러면 나머지 단계도 각자finally로 채널을 닫으며 끝난다. 마지막으로Task.WhenAll이 원래의 예외를RunAsync밖으로 던진다. 어느 단계도 무한히 기다리지 않고, 예상하지 못한 오류가 조용히 사라지지도 않는다. - 바깥 토큰을 받아 연결된
CancellationTokenSource를 쓴다.
호출하는 쪽에서는public async Task<PipelineResult> RunAsync( IReadOnlyList<string> input, int stopAfter, CancellationToken outer = default) { using var cts = CancellationTokenSource.CreateLinkedTokenSource(outer); // 나머지는 그대로 }Console.CancelKeyPress이벤트에서e.Cancel = true로 종료를 막고 별도의CancellationTokenSource의Cancel()을 부른 뒤, 그 토큰을RunAsync에 넘긴다. 연결된 소스는 바깥 토큰이나 내부Cancel()어느 쪽이든 취소되므로, 접수 단계의 자체 취소와 외부 취소가 같은 경로로 종료된다.