Half-Open Socket으로 인한 Airflow Stuck Task 디버깅
Airflow 운영 중 task가 queued 상태에서 약 20분간 멈춰있다가 자동 실패 처리되고, 좀비 Pod가 클러스터에 잔존하는 현상을 디버깅한 기록입니다. 표면적 현상만 보면 어디 한 곳이 명확하게 죽은 것 같지 않은데도 task가 계속 stuck되어 원인 추적이 오래 걸렸습니다.
문제 상황
CeleryExecutor + RabbitMQ 조합으로 Airflow를 운영 중이었고, KubernetesPodOperator(이하 KPO)를 주로 사용하고 있었습니다. 어느 날부터 다음과 같은 패턴이 간헐적으로 관찰되기 시작합니다.
03:52:21 task try=1 running ← KPO가 Pod를 만들고 base container까지 정상 종료
04:14:34 task try=1 FAILED ← 약 21분 후 자동 fail
04:24:43 task try=2 running ← scheduler가 자동 retry
04:25:56 task try=2 success ← 재시도는 정상 성공
특징을 정리하면 이렇습니다.
clear_number=0,user_actions=0→ 사용자가 UI에서 Clear/Mark Success를 누른 적이 없습니다.external_trigger=1이지만 외부 매니저 DAG의 일반 trigger 한 번뿐, cascade clear도 아닙니다.max_tries=1이지만 try=2가 발생합니다. 이는 scheduler의cleanup_stuck_queued_tasks가 호출한 자동 retry path였습니다.- try=1에서 만든 Pod의 base container는 정상적으로 exit 0을 반환했지만, KPO 입장에서는 Pod의 상태 변화 신호를 받지 못해 좀비 Pod로 남아 있습니다.
- stuck 지속 시간을 30일치 모아 정량 분석하니 mean 1283.9s (median 1288.4s), 80% 이상이 [1100s, 1400s] 구간에 분포 → deterministic timeout임이 명확합니다.
이 1283.9s는 정확히 task_queued_timeout=1200s + check_interval 1~2 cycle(120s)와 일치합니다. 즉 다른 어떤 layer도 dead detection을 하지 못하고, 가장 늦은 fallback인 scheduler의 stuck task cleanup만이 발화하고 있다는 뜻입니다.
등장 인물과 통신 경로
원인을 추적하기 위해 먼저 시스템에 누가 있는지 정리합니다.
flowchart LR
classDef ctrl fill:#dbeafe,stroke:#1d4ed8
classDef msg fill:#fef3c7,stroke:#a16207
classDef work fill:#fce7f3,stroke:#be185d
classDef exec fill:#dcfce7,stroke:#15803d
S["Scheduler<br/>(Airflow scheduler)"]:::ctrl
B["Broker<br/>(RabbitMQ)"]:::msg
W["Worker<br/>(Celery worker)"]:::work
P["Pod<br/>(K8s)"]:::exec
S -->|"publish task"| B
B -->|"deliver task"| W
W -->|"create pod"| P
여기서 한 가지 중요한 사실은, Worker가 broker하고만 통신하는 게 아니라는 점입니다.
flowchart LR
classDef worker fill:#fef3c7,stroke:#a16207,color:#78350f
classDef broker fill:#dcfce7,stroke:#15803d,color:#14532d
classDef db fill:#dbeafe,stroke:#1e40af,color:#1e3a8a
W["Worker<br/>(MainProcess + ForkPool)"]:::worker
B["Broker<br/>(RabbitMQ)"]:::broker
DB["metaDB<br/>(MySQL)"]:::db
B -->|"① task 메시지 deliver"| W
W -->|"② ACK frame (task 완료 후)"| B
W <-->|"③ heartbeat frame (120s 주기, 양방향)"| B
W -->|"④ celery_taskmeta INSERT (task 결과)"| DB
W -->|"⑤ task_instance UPDATE (실행 상태)"| DB
| 채널 | 방향 | broker 경유 |
|---|---|---|
| ① task deliver | Broker → Worker | O |
| ② ACK | Worker → Broker | O |
| ③ heartbeat frame | 양방향 (예: 120s 주기) | O |
④ celery_taskmeta INSERT |
Worker → DB | X (DB 직접) |
⑤ task_instance UPDATE |
Worker(task) → DB | X (DB 직접) |
scheduler는 task 상태 판단을 broker가 아닌 DB(metaDB) 만 보고 합니다. 따라서 broker-worker 간 메시지 전달이 silent하게 끊겨도, scheduler는 한참 뒤 (task_queued_timeout이 만료될 때) 비로소 이상을 감지할 수 있습니다.
Half-Open Socket
위 현상의 원인은 broker와 worker 사이의 TCP connection이 half-open 상태가 되었기 때문이었습니다.
Half-Open Socket: TCP connection의 양쪽 종단 중 한쪽은 비정상 종료/단절되었는데, 다른 한쪽은 여전히 정상으로 인식하고 있는 상태.
이게 가능한 이유는 TCP의 신뢰성 보장이 데이터를 send할 때만 동작하기 때문입니다. send가 없으면 상대방이 죽었어도 알 길이 없습니다.
sequenceDiagram
participant A as Worker
participant B as Broker
Note over A,B: 시나리오 1 — Worker가 send를 한다면
A->>B: 데이터 send
B-->>A: ACK
Note over A: Worker는 ACK로 broker alive 확인
Note over A,B: 시나리오 2 — Worker는 recv 대기만
Note over A: Worker는 그냥 기다림 (recv block)
B-xB: Broker쪽 socket이 죽음 (NAT timeout 등)
Note over B: FIN/RST도 못 보냄
Note over A: Worker는 영원히 기다림<br/>(send가 없어 ACK도 없고 알 길이 없음)
Celery worker는 task를 처리하지 않는 idle 구간 동안 recv()로 broker의 deliver를 기다립니다. 이 구간에서 worker는 broker에 데이터를 거의 보내지 않습니다(heartbeat는 별도 channel/timer로 발송됨). 그래서 특정 channel의 packet이 silent하게 drop되어도 worker application은 그 사실을 즉시 알 수 없습니다.
5개의 감시자, 그러나 모두 무력화
“TCP가 그렇다 쳐도, 우리에겐 여러 단계의 dead detection 메커니즘이 있지 않나?”하는 의문이 생깁니다. 실제로 Celery + RabbitMQ + Linux 조합에는 다음과 같은 timeout/heartbeat 메커니즘들이 있습니다.
flowchart TB
classDef weak fill:#fee2e2,stroke:#b91c1c
classDef mid fill:#fef3c7,stroke:#a16207
classDef strong fill:#dcfce7,stroke:#15803d
classDef trigger fill:#fce7f3,stroke:#be185d
subgraph app["Application layer (py-amqp)"]
AA["read_timeout = None<br/>= recv() 무한 block"]:::weak
AB["broker_heartbeat = 120s<br/>(양방향 약속)"]:::mid
AC["SO_KEEPALIVE = 1"]:::strong
end
subgraph os["OS Kernel layer"]
OA["tcp_keepalive total = 1185s<br/>(SO_KEEPALIVE 켜져있어야 fire)"]:::mid
OD["tcp_retries2 = 15 → ≈924s<br/>(unacked SENT data가 있어야 fire)"]:::strong
end
subgraph broker["Broker (RabbitMQ < 3.8.15)"]
BA["AMQP heartbeat = 120s"]:::mid
BB["consumer_timeout (없음)"]:::weak
end
subgraph sched["Airflow Scheduler"]
SA["task_queued_timeout = 1200s"]:::trigger
SB["check_interval = 120s"]:::trigger
end
그런데 위에서 본 정량 데이터(stuck duration ≈ 1283s)는 이 중 scheduler의 1200s task_queued_timeout만이 fire되었음을 의미합니다. 나머지는 왜 동작하지 않았을까요? 각각 발화 조건이 만족되지 않습니다.
| Layer | 발화 안 한 이유 |
|---|---|
py-amqp read_timeout |
기본값 None. SO_RCVTIMEO 자체가 설정 안 되어 recv()가 무한 block. timeout이라는 개념 자체가 없음. |
| AMQP heartbeat (120s) | py-amqp의 Hub event loop가 다른 channel event로 정상 깨어나며 heartbeat timer callback을 정상 발화. broker는 client를 alive로 봄. |
| TCP keepalive (1185s) | SO_KEEPALIVE는 켜져 있고 probe도 발사하지만, broker의 OS가 정상 ACK 응답을 함. keepalive 입장에선 connection 정상. |
TCP tcp_retries2 (924s) |
이 timer는 unacked SENT data에 대해서만 의미가 있음. worker가 idle하게 recv만 하는 동안 send가 없으므로 발화 조건 미충족. |
RMQ consumer_timeout |
사용 중이던 broker 버전에는 아예 존재하지 않는 기능 (3.8.15+에서 추가). broker self-heal 불가. |
Airflow task_queued_timeout |
유일하게 발화. metaDB의 queued_dttm 컬럼만 비교하므로 네트워크 상태와 무관. |
핵심은 “send가 없으면 OS layer는 아무것도 못 한다” 입니다. consumer가 idle하게 recv만 대기하는 동안은 OS 입장에서 진행 중인 통신이 없는 것과 같기 때문에, tcp_retries2도 keepalive의 dead 판정도 작동하지 않습니다. 그리고 application layer에서 깨워줄 read_timeout도 비활성화되어 있으니, half-open이 발생한 채로 무한히 sleep합니다.
py-amqp의 read_timeout=None
py-amqp는 Celery가 RabbitMQ와 통신할 때 사용하는 AMQP 라이브러리입니다. 코드를 따라가면 connection 생성 시 read_timeout 인자가 기본 None으로 들어가고, 이 값이 None인 경우 SO_RCVTIMEO socket option이 설정되지 않습니다.
sequenceDiagram
autonumber
participant App as Worker py-amqp
participant OS as Worker OS kernel
participant B as Broker
App->>OS: socket.recv() 호출 + (read_timeout=None)
Note over OS: SO_RCVTIMEO 미설정 → 무한 대기
Note over App: sleep
B-xOS: 메시지를 보내려 하지만 packet drop
Note over OS: socket buffer에 데이터 안 들어옴
Note over App: 영원히 sleep<br/>scheduler가 1200s 후 fail 처리할 때까지
만약 read_timeout이 예를 들어 60초로 설정되어 있다면 어떨까요?
sequenceDiagram
autonumber
participant App as Worker py-amqp
participant OS as Worker OS kernel
participant B as Broker
App->>OS: socket.recv() 호출 + (read_timeout=60)
Note over OS: SO_RCVTIMEO=60 설정
B-xOS: packet drop
Note over OS: 60s 동안 buffer empty
OS->>App: wake up + EAGAIN
App->>App: heartbeat_check() / connection state 재검증
App->>App: half-open 감지 → 재연결
Note over App,B: broker가 unacked 메시지 redeliver
Note over App: 1분 안에 자동 복구
recv()가 60초마다 EAGAIN으로 깨어나면 py-amqp는 connection 상태를 재검증할 기회를 얻습니다. 이때 heartbeat가 실제로 왔는지 확인하고, 비정상 상태가 감지되면 connection을 닫고 재연결합니다. 그러면 broker는 unacked 메시지를 다른 정상 consumer에게 redeliver하게 되고, 1분 안에 자연스럽게 복구됩니다.
두 timeout의 stack 관계
여기서 한 가지 헷갈리는 부분을 정리할 필요가 있습니다. “AMQP heartbeat이 어차피 120s마다 주고받는데, read_timeout이 따로 필요한가?”
답은 read_timeout이 없으면 heartbeat 체크도 무력화된다는 것입니다. heartbeat frame도 결국 TCP packet이고, application이 그것을 보려면 OS가 application을 깨워야 합니다. application이 영원히 sleep 중이면 heartbeat frame이 socket buffer에 도착해도 처리할 코드가 실행되지 않습니다.
flowchart LR
classDef l1 fill:#dbeafe,stroke:#1e40af
classDef l2 fill:#fef3c7,stroke:#a16207
T2["read_timeout (OS 시계)<br/>application을 깨우는 알람"]:::l1
Wake["application 깨어남"]
T3["heartbeat check (AMQP 시계)<br/>heartbeat 수신 여부 검증"]:::l2
T2 -->|"X초마다 깨움"| Wake
Wake -->|"깨어나야 체크 가능"| T3
| 관계 | 내용 |
|---|---|
| heartbeat 체크는 read_timeout 위에 stack | heartbeat frame도 socket으로만 옴. read_timeout이 깨워야 heartbeat 검증 코드가 실행됨. |
| read_timeout이 None이면 heartbeat 무력 | application이 영원히 sleep → heartbeat 검증 코드 실행 안 됨 → broker 죽어도 모름. |
| read_timeout만으로는 부족 | “데이터 안 옴”만 알지 “broker 죽음”까지 판정 못 함. heartbeat 검증 로직과 결합되어야 의미를 가짐. |
흔히 “AMQP heartbeat이 있으니 dead detection은 broker가 알아서 해주겠지”라고 생각하기 쉽지만, 위 stack 구조 때문에 application의 read_timeout이 끊어주는 깨움 신호가 없으면 heartbeat 체크 자체가 동작하지 않습니다.
또한, heartbeat은 보통 단일 connection 단위로 동작하고, Hub event loop의 timer는 다른 channel event로도 깨어날 수 있습니다. 즉 “문제가 생긴 채널은 silent drop인데, heartbeat은 다른 채널을 타고 정상 송수신되어 broker 입장에서는 client alive로 보이는 상황” 이 가능합니다. 이때 broker는 의심할 이유가 없으니 절대 먼저 끊지 않습니다.
해결책
가장 적은 변경으로 가장 넓은 케이스를 cover하는 방법은 application layer에서 직접 감지하도록 만드는 것입니다. Celery의 broker transport options를 통해 socket timeout을 명시적으로 설정합니다.
AIRFLOW__CELERY_BROKER_TRANSPORT_OPTIONS__SOCKET_TIMEOUT=60
AIRFLOW__CELERY_BROKER_TRANSPORT_OPTIONS__SOCKET_CONNECT_TIMEOUT=10
AIRFLOW__SCHEDULER__TASK_QUEUED_TIMEOUT=1200
SOCKET_TIMEOUT=60: py-amqp의read_timeout기본값을 변경. recv()가 60초마다 EAGAIN으로 깨어나 connection을 재검증.SOCKET_CONNECT_TIMEOUT=10: 초기 connection 단계의 무한 hang 방지.TASK_QUEUED_TIMEOUT=1200: scheduler stuck task fallback. 최후의 방어선으로 유지.
위 설정을 한쪽 환경에만 먼저 적용하고 30일 운영 데이터를 다른 환경(설정 미적용)과 A/B 비교해 봤습니다.
| 환경 | SOCKET_TIMEOUT | 30일 task 처리량 | REVOKED celery_taskmeta |
|---|---|---|---|
| 미적용 | None (default) | ≈ 1.5M | 23건 (≈ 0.00153%) |
| 적용 | 60s | ≈ 200K | 0건 |
미적용 환경 발생률을 그대로 적용해 보면 200K 처리량에서 통계상 약 3건이 발생해야 하지만 실제로는 0건이었습니다. 30일 운영으로 정량 확정된 효과입니다.
다른 fix 후보들
위 application layer 설정 외에도 이론적으로 가능한 fix들이 있습니다.
graph TD
classDef fix fill:#dcfce7,stroke:#15803d
classDef cur fill:#fee2e2,stroke:#b91c1c
H["Half-Open Socket 발생<br/>(NAT timeout / 네트워크 glitch)"]:::cur
F1["read_timeout 설정<br/>(application layer)"]:::fix
F2["RMQ 3.8.15+ 업그레이드<br/>(consumer_timeout 추가)"]:::fix
F3["TCP_USER_TIMEOUT 명시 설정<br/>(kernel layer)"]:::fix
F4["broker_heartbeat 짧게 (예: 30s)"]:::fix
H --> F1
H --> F2
H --> F3
H --> F4
F1 -->|"가장 넓은 cover"| OK["좀비 안 생김"]
F2 -->|"broker side self-heal"| OK
F3 -->|"send가 있어야 fire"| Partial["부분 효과"]
F4 -->|"Hub timer 정상 시 무력"| Partial
- RMQ 3.8.15+ 업그레이드: broker에
consumer_timeout이 추가되어, broker가 일정 시간 ack 받지 못한 consumer를 끊고 메시지를 redeliver합니다. broker side self-heal이라는 점에서 가장 깔끔하지만 broker upgrade 비용이 큽니다. - TCP_USER_TIMEOUT: kernel layer에서 unacked SENT data의 누적 시간을 직접 제한. 강력하지만 idle consumer처럼 send가 없는 케이스에선 부분적 효과만 있습니다.
- broker_heartbeat을 짧게: 빠른 dead detection이 가능해 보이지만, 위에서 본 것처럼 Hub timer가 정상 작동하는 상황에서는 broker가 client를 alive로 봐서 무력화됩니다.
application layer의 read_timeout 설정이 가장 적은 변경으로 가장 넓은 범위의 half-open 케이스를 cover합니다.
정리
이 디버깅에서 배운 것을 정리하면 다음과 같습니다.
- TCP는 send할 때만 신뢰성을 보장한다. recv만 하는 idle consumer에 대해서는 dead detection이 자동으로 일어나지 않습니다.
- 여러 layer의 timeout이 있어도, 발화 조건을 만족하지 못하면 모두 무력하다. 특히 application layer의 read_timeout이 없으면 heartbeat 검증조차 실행되지 않습니다.
- library의 기본값을 의심해야 한다. py-amqp의
read_timeout=None은 정상 환경에서는 아무 문제가 없지만, 네트워크 glitch가 발생한 순간 가장 약한 고리가 됩니다. - stuck duration의 분포가 결정적 timeout과 일치한다면, 원인은 거의 확정된다. mean 1283.9s = scheduler의 1200s + check_interval 1~2 cycle. 이 시간 일치는 다른 어떤 layer도 발화하지 않고 scheduler fallback만 작동했음을 말해줍니다.
- DB 기반 scheduler는 broker 상태와 독립적이다. broker와 worker의 통신이 silent하게 끊겨도 scheduler는 metaDB의
queued_dttm만 비교해 fail 처리할 수 있고, 결국 이 fallback이 시스템을 살립니다. 다만 발화가 늦어 좀비 Pod가 잔존합니다.
운영 환경의 안정성은 자주 발화하지 않는 timeout일수록 더 주의 깊게 살펴봐야 한다는 교훈을 남긴 디버깅이었습니다.