# Process Stream - 스트림 데이터 처리
## 개념 설명
실시간 스트림 데이터를 비동기적으로 처리하는 핵심 개념입니다. 대용량 데이터 스트림을 효율적으로 처리하면서 백프레셔(backpressure) 제어, 동시성 관리, 에러 복구 등의 고급 기능을 제공합니다.
## 핵심 특징
### + 백프레셔 처리
- 처리 속도보다 입력이 빠를 때 자동으로 흐름 제어
- 메모리 사용량을 일정 수준으로 유지
- Channel(C#) / Queue(Python) 기반 버퍼링
### ⚡ 동시성 제어
- 설정 가능한 최대 동시 처리 개수
- SemaphoreSlim(C#) / asyncio.Semaphore(Python)로 리소스 관리
- CPU 코어 수에 따른 자동 최적화
### ?+ 에러 복구
- 개별 아이템 처리 실패가 전체 스트림을 중단시키지 않음
- 에러 콜백을 통한 커스텀 에러 처리
- Continue-on-error 옵션으로 유연한 에러 정책
### + 조합 가능성
- 여러 변환 단계를 체인으로 연결
- 함수형 프로그래밍 스타일 지원
- 필터링, 매핑 등 유틸리티 함수 제공
## 인터페이스 설계
### C# 인터페이스 (BSD 스타일)
```csharp
public interface IStreamProcessor<TInput, TOutput>
{
// 기본 옵션으로 처리
IAsyncEnumerable<TOutput> ProcessAsync(
IAsyncEnumerable<TInput> inputStream,
Func<TInput, Task<TOutput>> transform,
CancellationToken cancellationToken = default);
// 커스텀 옵션으로 처리
IAsyncEnumerable<TOutput> ProcessAsync(
IAsyncEnumerable<TInput> inputStream,
Func<TInput, Task<TOutput>> transform,
StreamProcessorOptions<TInput> options,
CancellationToken cancellationToken = default);
}
```
### Python 인터페이스
```python
class StreamProcessor(ABC):
@abstractmethod
async def process(
self,
input_stream: AsyncIterator[T],
transform: Callable[[T], Awaitable[U]],
options: Optional[StreamProcessorOptions] = None
) -> AsyncIterator[U]:
pass
```
## 구현 세부사항
### C# 구현 특징
- **Channel<T>**: 백프레셔를 위한 bounded channel 사용
- **SemaphoreSlim**: 동시성 제어
- **Task.WhenAll**: 모든 처리 작업 완료 대기
- **IAsyncEnumerable**: 지연 실행과 메모리 효율성
- **BSD 스타일**: 가독성을 위한 중괄호 새 줄 배치
### Python 구현 특징
- **asyncio.Queue**: 백프레셔를 위한 maxsize 제한 큐
- **asyncio.Semaphore**: 동시성 제어
- **asyncio.gather**: 병렬 작업 관리
- **AsyncIterator**: 지연 실행과 메모리 효율성
- **Type Hints**: 타입 안전성 보장
## 사용 예시
### 기본 데이터 변환
```csharp
// C# 예시
var processor = new StreamProcessor<string, int>();
var numbers = processor.ProcessAsync(
textStream,
async text =>
{
await Task.Delay(10); // 처리 시뮬레이션
return int.Parse(text);
}
);
await foreach (var number in numbers)
{
Console.WriteLine($"Parsed: {number}");
}
```
```python
# Python 예시
processor = AsyncStreamProcessor()
async def parse_number(text: str) -> int:
await asyncio.sleep(0.01) # 처리 시뮬레이션
return int(text)
async for number in processor.process(text_stream, parse_number):
print(f"Parsed: {number}")
```
### 에러 처리가 포함된 처리
```csharp
// C# 에러 처리
var options = new StreamProcessorOptions<string>
{
MaxConcurrency = 4,
Continue = true,
= (ex, input) => Console.WriteLine($"Failed to process {input}: {ex.Message}")
};
await foreach (var result in processor.ProcessAsync(dataStream, transform, options))
{
Console.WriteLine($"Success: {result}");
}
```
```python
# Python 에러 처리
def error_handler(ex: Exception, item: str):
print(f"Failed to process {item}: {ex}")
options = StreamProcessorOptions(
max_concurrency=4,
continue_on_error=True,
on_error=error_handler
)
async for result in processor.process(data_stream, transform, options):
print(f"Success: {result}")
```
### 스트림 체이닝
```csharp
// C# 체이닝
var processor1 = new StreamProcessor<string, int>();
var processor2 = new StreamProcessor<int, string>();
var result = processor2.ProcessAsync(
processor1.ProcessAsync(stringStream, ParseInt),
async num => $"Number: {num * 2}"
);
```
```python
# Python 체이닝 (유틸리티 함수 사용)
filtered = filter_stream(raw_stream, lambda x: x > 0)
squared = map_stream(filtered, lambda x: x * x)
async for result in squared:
print(f"Filtered and squared: {result}")
```
## 성능 특성
### 메모리 사용량
- **O(BufferSize)**: 설정된 버퍼 크기에 비례한 일정한 메모리 사용
- **스트리밍 처리**: 전체 데이터를 메모리에 로드하지 않음
- **백프레셔**: 메모리 부족 방지를 위한 자동 흐름 제어
### 처리 성능
- **병렬 처리**: MaxConcurrency 설정으로 처리량 조절
- **비동기 I/O**: I/O 바운드 작업에 최적화
- **지연 실행**: 필요할 때만 데이터 처리
### 확장성
- **수평 확장**: 여러 인스턴스로 분산 처리 가능
- **수직 확장**: 동시성 수준 조정으로 리소스 활용 최적화
## 적용 사례
### 실시간 데이터 처리
- 로그 스트림 분석
- 센서 데이터 처리
- 실시간 메트릭 수집
### ETL 파이프라인
- 대용량 데이터 변환
- 데이터 정제 및 검증
- 포맷 변환
### 이벤트 처리
- 메시지 큐 처리
- 이벤트 스트림 변환
- 실시간 알림 시스템
## 관련 개념
- **transform-batch**: 배치 단위 처리가 필요한 경우
- **handle-events**: 이벤트 기반 처리와 조합
- **validate-input**: 입력 검증과 함께 사용
- **cache-data**: 처리 결과 캐싱
- **retry-operations**: 실패한 처리 재시도
어떠냐
이 댓글은 게시물 작성자가 삭제하였습니다.
면접은 잘하겠노 실무경험을 물어보겠지만
나는 걍 프리랜서가 나아
이 댓글은 게시물 작성자가 삭제하였습니다.
오 어떻게해 ☆☆☆☆☆☆☆☆☆☆☆☆☆☆☆