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 서버·임시 포트·구독자는 검사기가 준비하고 종료합니다. 이 실습은 커널 버퍼 포화·인터넷 성능·영속 전달·실제 재접속 복구를 보장하지 않습니다. 실습 세션이 끝나면 파일은 유지되지 않으니 필요한 코드는 종료 전에 별도로 보관하세요.

끝없이 쌓이는 큐에 상한을 건다

delivery.py에 make_queue(capacity)를 구현하세요. 비어 있는 asyncio.Queue를 반환하고 maxsize는 입력과 같아야 합니다. 양의 int만 허용하며 bool·실수·문자열·0·음수는 ValueError입니다.

asyncio.Queue에서 0은 무제한입니다. type 검사와 범위 검사를 나누고 표준 라이브러리를 사용하세요.

기다리되 취소된 발행은 남기지 않는다

async offer_wait(queue, item)를 추가하세요. 공간이 생길 때까지 기다려 넣고 None을 반환합니다. 대기 중 취소는 CancelledError로 전달하며 기존 항목과 순서를 보존합니다. 취소된 항목을 나중에 삽입하지 않습니다.

put_nowait와 await put의 포화 동작은 다릅니다. 취소를 정상 성공으로 바꾸지 마세요.

옛 좌표를 버리고 현재 위치를 남긴다

동기 함수 offer_latest(queue, item)를 추가하세요. 공간이 있으면 삽입 후 None, 가득 차면 가장 오래 기다린 항목 하나를 제거하고 새 항목을 넣은 뒤 제거한 항목을 반환합니다. 버린 항목의 task_done도 한 번 대응시킵니다.

용량 2에 10,11이 있을 때 12가 오면 11,12가 남습니다. 처리 중 항목은 이 큐에서 제거할 수 없습니다.

멈춘 구독자를 정책 예외로 알린다

Exception 하위 클래스 SlowConsumer와 동기 offer_disconnect(queue, item)를 구현하세요. 공간이 있으면 삽입 후 None을 반환합니다. 가득 차면 기존 큐를 바꾸지 않고 SlowConsumer를 냅니다. 연결을 닫는 일은 이 함수의 책임이 아닙니다.

한 함수는 포화 정책을 판단하고 연결 소유자는 회수합니다. 조용히 버리거나 기다리는 것은 다른 정책입니다.

중복은 건너뛰고 누락은 숨기지 않는다

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입니다. 메모리 커서이며 영속 업무 처리는 구현하지 않습니다.

중복 검사와 갭 검사의 순서가 중요합니다. 7 다음에 9가 오면 8을 처리한 근거가 없습니다.

송신과 업무 확인에 한 번만 시간을 준다

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를 닫지 않고 오류·취소를 전달합니다.

asyncio.wait_for에는 매번 남은 시간을 줍니다. 테스트는 시계를 주입하며 drain 0.8초와 ACK 0.3초가 전체 1초를 넘는지 확인합니다.

비어 있거나 실패해도 연결을 회수한다

async serve_queue(reader, writer, queue, timeout)를 추가하세요. 양의 유한 timeout이 전달됩니다. 반복해서 get한 seq를 send_one으로 처리하며, 그 호출이 끝나거나 실패·취소된 뒤에만 인수한 항목에 task_done을 정확히 한 번 호출합니다. 빈 큐 대기를 포함해 모든 종료 경로에서 writer.close 후 wait_closed를 timeout 안에 기다립니다. 종료 대기의 시간 초과·ConnectionError에는 writer.transport.abort로 회수하고 원래 송신 오류·취소를 숨기지 않습니다. 아직 큐에 남은 항목 정리는 호출자의 책임입니다.

항목 소유권의 try/finally와 연결 소유권의 try/finally는 범위가 다릅니다. ACK 전에 task_done을 부르면 join이 먼저 풀립니다.

느린 한 명을 분리하고 생중계를 계속한다

동기 broadcast(queues, item, disconnect)를 완성하세요. queues는 이름→큐 딕셔너리이고 순회 스냅샷의 각 큐에 offer_disconnect를 적용합니다. SlowConsumer인 대상만 딕셔너리에서 제거하고 disconnect(name)를 한 번 호출한 뒤 다른 구독자에게 계속 전달합니다. 다른 예외는 전달하고 정상 반환은 None입니다. 콜백은 예외 없이 해당 작업자 취소·연결 회수·잔여 큐 정리를 담당합니다. 실제 TCP 검사에서 slow는 첫 ACK를 보류하고 fast는 매번 확인합니다. 0~8 이벤트를 fast가 모두 받고 slow만 분리되어야 합니다.

첫 분리 뒤 return하지 마세요. 검사기가 실제 서버와 두 클라이언트를 띄우므로 직접 서버를 상시 실행할 필요는 없습니다.