LabHub
배우기 러닝패스 코스

느린 구독자가 생중계를 멈췄다 · 확인·기한·회수의 주인을 정한다 · 실습

느린 구독자를 격리해 생중계를 구한다

LabHub 에서 이어서 보기

목표

큐 정책과 ACK 기한, 연결 수명을 구현해 실제 TCP에서 느린 구독자를 격리합니다.

왜 중요한가

구독자 한 명의 지연이 전체 발행을 막거나 메모리를 계속 사용하게 할 수 있습니다. 데이터의 손실 허용 범위와 업무 확인을 구분하고, 취소 때도 자원을 회수하는 프로그램을 만듭니다. Python 함수·예외·async/await 기초가 필요합니다. 표준 라이브러리만 사용하며 설치나 인터넷은 필요 없습니다.

단계

1. 끝없이 쌓이는 큐에 상한을 건다 — delivery.py에 make_queue(capacity)를 구현하세요. 비어 있는 asyncio.Queue를 반환하고 maxsize는 입력과 같아야 합니다. 양의 int만 허용하며 bool·실수·문자열·0·음수는 ValueError입니다.
2. 기다리되 취소된 발행은 남기지 않는다 — async offer_wait(queue, item)를 추가하세요. 공간이 생길 때까지 기다려 넣고 None을 반환합니다. 대기 중 취소는 CancelledError로 전달하며 기존 항목과 순서를 보존합니다. 취소된 항목을 나중에 삽입하지 않습니다.
3. 옛 좌표를 버리고 현재 위치를 남긴다 — 동기 함수 offer_latest(queue, item)를 추가하세요. 공간이 있으면 삽입 후 None, 가득 차면 가장 오래 기다린 항목 하나를 제거하고 새 항목을 넣은 뒤 제거한 항목을 반환합니다. 버린 항목의 task_done도 한 번 대응시킵니다.
4. 멈춘 구독자를 정책 예외로 알린다 — Exception 하위 클래스 SlowConsumer와 동기 offer_disconnect(queue, item)를 구현하세요. 공간이 있으면 삽입 후 None을 반환합니다. 가득 차면 기존 큐를 바꾸지 않고 SlowConsumer를 냅니다. 연결을 닫는 일은 이 함수의 책임이 아닙니다.
5. 중복은 건너뛰고 누락은 숨기지 않는다 — Exception 하위 클래스 Gap과 Cursor(last=-1)를 추가하세요. last는 인스턴스별 공개 속성이고 -1~2147483647 범위 int입니다. accept(seq)는 0~2147483647 int만 받습니다. seq<=last면 False와 상태 유지, seq==last+1이면 last 갱신 후 True, 더 큰 순번이면 Gap과 상태 유지입니다. bool 등 잘못된 타입·범위는 ValueError입니다. 메모리 커서이며 영속 업무 처리는 구현하지 않습니다.
6. 송신과 업무 확인에 한 번만 시간을 준다 — async send_one(reader, writer, seq, timeout, clock=None)을 만드세요. seq는 5단계와 같은 유효 순번이며 timeout은 bool을 제외한 양의 유한 int/float이고 잘못되면 ValueError입니다. 기본 시계는 time.monotonic, 주입 시계는 인자 없이 초를 반환합니다. ASCII EVENT 순번 뒤 개행을 writer.write로 한 번 쓰고 drain을 기다린 다음 reader.readline으로 정확히 같은 순번의 ACK와 개행을 받아야 None입니다. EOF·잘못된 ACK는 ConnectionError, 전체 기한 초과는 TimeoutError입니다. drain과 ACK는 같은 마감 시각을 공유하며 완료 후에도 기한을 확인합니다. 이 함수는 writer를 닫지 않고 오류·취소를 전달합니다.
7. 비어 있거나 실패해도 연결을 회수한다 — async serve_queue(reader, writer, queue, timeout)를 추가하세요. 양의 유한 timeout이 전달됩니다. 반복해서 get한 seq를 send_one으로 처리하며, 그 호출이 끝나거나 실패·취소된 뒤에만 인수한 항목에 task_done을 정확히 한 번 호출합니다. 빈 큐 대기를 포함해 모든 종료 경로에서 writer.close 후 wait_closed를 timeout 안에 기다립니다. 종료 대기의 시간 초과·ConnectionError에는 writer.transport.abort로 회수하고 원래 송신 오류·취소를 숨기지 않습니다. 아직 큐에 남은 항목 정리는 호출자의 책임입니다.
8. 느린 한 명을 분리하고 생중계를 계속한다 — 동기 broadcast(queues, item, disconnect)를 완성하세요. queues는 이름→큐 딕셔너리이고 순회 스냅샷의 각 큐에 offer_disconnect를 적용합니다. SlowConsumer인 대상만 딕셔너리에서 제거하고 disconnect(name)를 한 번 호출한 뒤 다른 구독자에게 계속 전달합니다. 다른 예외는 전달하고 정상 반환은 None입니다. 콜백은 예외 없이 해당 작업자 취소·연결 회수·잔여 큐 정리를 담당합니다. 실제 TCP 검사에서 slow는 첫 ACK를 보류하고 fast는 매번 확인합니다. 0~8 이벤트를 fast가 모두 받고 slow만 분리되어야 합니다.

참고

모든 함수는 /root/realtime/delivery.py 하나에 둡니다. mkdir -p /root/realtime로 작업 폴더를 만드세요. 예시의 아직 구현하지 않은 함수는 틀로 두되 완료한 함수를 덮어쓰지 마세요. 채점은 이전 단계 계약도 검사합니다. 한 단계 실행은 5초 제한이며 학습자가 작성할 시간의 제한이 아닙니다. TCP 서버·임시 포트·구독자는 검사기가 준비하고 종료합니다. 이 실습은 커널 버퍼 포화·인터넷 성능·영속 전달·실제 재접속 복구를 보장하지 않습니다. 실습 세션이 끝나면 파일은 유지되지 않으니 필요한 코드는 종료 전에 별도로 보관하세요.

단계 8개

  1. 끝없이 쌓이는 큐에 상한을 건다
  2. 기다리되 취소된 발행은 남기지 않는다
  3. 옛 좌표를 버리고 현재 위치를 남긴다
  4. 멈춘 구독자를 정책 예외로 알린다
  5. 중복은 건너뛰고 누락은 숨기지 않는다
  6. 송신과 업무 확인에 한 번만 시간을 준다
  7. 비어 있거나 실패해도 연결을 회수한다
  8. 느린 한 명을 분리하고 생중계를 계속한다