Airflow에서 Kestra로 넘어가며 겪은 파이프라인 전환 기록
TuBrief Editorial
July 1, 2026
0
Computing/SoftwareWritten with AI assistance from the source video. The video is the authority.
More from the community
Comments (0)
Log in to leave a comment
No posts yet
Written with AI assistance from the source video. The video is the authority.
Log in to leave a comment
No posts yet
오케스트레이션 정의와 비즈니스 로직이 파이썬 코드로 떡칠된 Apache Airflow 환경은 엔지니어링 리소스가 부족한 스타트업에게 재앙입니다. 관리에 손이 너무 많이 갑니다. 서비스 중단 없이 안전하게 이전을 마치려면 기존 파이프라인을 쪼개어 옮기는 스트랭글러 피그 패턴(Strangler Fig Pattern)을 써야 합니다.
Airflow TaskFlow API 내부 로직을 Kestra의 선언적 YAML 스크립트 태스크로 매핑하는 과정은 생각보다 단순합니다.
containerImage: python:3.11-slim 기반의 Docker 러너로 지정합니다.Kestra.outputs({"key": value}) JSON 호출로 반환하고, 대용량 데이터프레임은 로컬 디스크에 CSV 파일로 저장합니다. 이 파일을 outputFiles 명세에 등록하면 Kestra 내장 저장소로 들어갑니다.inputFiles 속성에 상위 태스크의 파일 산출물 경로인 {{ outputs.task_id.outputFiles['filename.csv'] }}를 매핑해 격리된 컨테이너 파일 시스템으로 로드합니다.이 작업을 마치면 파이썬 패키지 의존성 충돌이 사라집니다. 데이터 파이프라인 전체 유지보수 시간도 눈에 띄게 줄어듭니다.
id: legacy_airflow_python_migration
namespace: production.data.migration
description: Airflow TaskFlow API 기반 로직을 Kestra Script 태스크로 분리 매핑한 마이그레이션 템플릿
tasks:
- id: python_extract_transform
type: io.kestra.plugin.scripts.python.Script
containerImage: python:3.11-slim
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
beforeCommands:
- pip install pandas kestra
outputFiles:
- "transformed_data.csv"
script: |
import pandas as pd
from kestra import Kestra
raw_records = [
{"transaction_id": "TX901", "amount": 150.50, "status": "completed"},
{"transaction_id": "TX902", "amount": 420.00, "status": "pending"}
]
df = pd.DataFrame(raw_records)
df_completed = df[df['status'] == 'completed']
df_completed.to_csv("transformed_data.csv", index=False)
Kestra.outputs({
"total_amount": float(df_completed['amount'].sum()),
"record_count": int(len(df_completed))
})
- id: write_to_warehouse
type: io.kestra.plugin.scripts.python.Script
containerImage: python:3.11-slim
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
inputFiles:
source_data.csv: "{{ outputs.python_extract_transform.outputFiles['transformed_data.csv'] }}"
beforeCommands:
- pip install pandas
script: |
import pandas as pd
refined_df = pd.read_csv("source_data.csv")
print(f"적재 대상 행 수: {len(refined_df)}")
개발자가 웹 콘솔 UI에서 YAML 코드를 직접 고치기 시작하면 형상 관리가 망가집니다. 운영 환경이 언제 터질지 모르는 시한폭탄이 됩니다. 이를 막으려면 운영 환경으로 들어오는 모든 플로우에 system.readOnly: true 라벨을 넣어 원격 UI 에디터를 잠가야 합니다. GitHub Actions와 Kestra CLI 도구인 kestractl을 엮어 배포를 자동화하는 방법입니다.
.github/workflows/deploy.yml 파일을 만들고 on: push와 pull_request 트리거를 겁니다.kestra-io/github-actions/validate-flows 액션으로 YAML 스키마 구조를 먼저 검증합니다. 원격 API에 영향을 주지 않는 드라이런 방식입니다.deploy-namespace-files 액션으로 외부 파이썬이나 SQL 리소스를 먼저 올립니다. 그 다음 deploy-flows 액션에 override: true 옵션을 주어 플로우를 최종 업로드합니다.의도치 않은 인프라 코드가 실시간 운영 환경에 그대로 반영되는 사고를 배포 단계에서 걸러낼 수 있습니다.
name: "Kestra Workflows Git Sync & Immutable Deploy"
on:
push:
branches: [ "main" ]
paths: [ "kestra/**" ]
pull_request:
branches: [ "main" ]
paths: [ "kestra/**" ]
jobs:
validate-flows:
runs-on: ubuntu-latest
steps:
- name: "코드 레포지토리 체크아웃"
uses: actions/checkout@v4
- name: "Kestra 플로우 구조 드라이런 검사"
uses: kestra-io/github-actions/validate-flows@main
with:
directory: ./kestra/flows
server: ${{ secrets.KESTRA_PROD_SERVER_URL }}
apiToken: ${{ secrets.KESTRA_API_TOKEN }}
deploy-workflows:
needs: validate-flows
if: github.event_name == 'push' && github.ref == 'refs/heads/main'
runs-on: ubuntu-latest
steps:
- name: "코드 레포지토리 체크아웃"
uses: actions/checkout@v4
- name: "네임스페이스 의존 파일 일괄 배포"
uses: kestra-io/github-actions/deploy-namespace-files@main
with:
directory: ./kestra/namespace-files
namespace: production.data
server: ${{ secrets.KESTRA_PROD_SERVER_URL }}
apiToken: ${{ secrets.KESTRA_API_TOKEN }}
- name: "Kestra 최종 워크플로우 파일 업로드"
uses: kestra-io/github-actions/deploy-flows@main
with:
directory: ./kestra/flows
server: ${{ secrets.KESTRA_PROD_SERVER_URL }}
apiToken: ${{ secrets.KESTRA_API_TOKEN }}
override: true
특정 도메인의 무거운 연산 워크플로우가 자원을 독점하면 결제나 인증 상태 체크 같은 핵심 파이프라인이 멈춥니다. 단일 인스턴스 안에서 비즈니스 도메인별로 네임스페이스 격리 레벨을 두고 자원 점유율을 묶어야 합니다. 가용 CPU와 메모리의 물리 상한선 이 시스템 총 수용량 R_{\text{worker_capacity}}를 넘지 않도록 제어하는 아키텍처 수식은 다음과 같습니다.
네임스페이스별로 concurrencyLimit 정책을 세우고 실행 환경 컨테이너 내부에 물리 리소스 상한을 지정합니다.
company.data.ingest는 동시 실행 제한을 5로 잡고 cpus: 1.0, memory: 1024Mb 자원을 줍니다.company.data.mart는 자원 독점을 막기 위해 동시 실행 제한을 2로 낮춥니다. 대신 cpus: 4.0, memory: 8192Mb로 연산 범위를 넓혀줍니다..env_encoded 파일을 도커 컴포즈의 env_file로 연동합니다. 플로우 YAML 안에서는 {{ secret('DATABASE_PROD_PASSWORD') }} 구문만 쓰도록 강제해서 네임스페이스끼리 보안 설정을 훔쳐보지 못하게 막습니다.데이터가 갑자기 몰려 특정 마트 생성 태스크가 넘치더라도 해당 격리 컨테이너 안에서 OOM(Out Of Memory) Kill 처리됩니다. 제어 엔진 전체가 같이 죽는 불상사를 피할 수 있습니다.
장애 원인도 모른 채 무작정 자동 재시도를 반복하면 스타트업의 지갑이 털립니다. 서드파티 API 쿼터만 고갈되어 2차 장애로 이어지기 십상입니다. 일시적인 네트워크 결함과 데이터 자체의 스키마 에러를 분리해서 대응해야 합니다.
| 예외 유형 | 장애 원인 예시 | 대응 액션 및 복구 전략 | 알림 제어 방식 |
|---|---|---|---|
| 일시적 네트워크 예외 | 외부 API 순단 (HTTP 503), DB 커넥션 병목 | retry 블록 기반 지수 백오프 가동 () | 최대 임계치까지 재시도 실패 시에만 장애 채널 알림 |
| 구조적 데이터 결함 | 정합성 불일치 (Null Value 유입), 스키마 위반 | allowFailed: true 우회 처리 또는 Fail 즉시 전파 |
StackTrace 분석 후 실시간 On-call 핫라인 전송 |
| 비핵심 태스크 장애 | 마케팅 메일 발송 지연, 단순 통계 누락 | allowWarning: true 지정으로 후속 작업 진행 |
경고 알림 전송 후 워크플로우는 정상 완료 처리 |
태스크 실패가 터지면 전역 에러 제어기를 돌려 Kestra 내장 함수 errorLogs()로 추출한 StackTrace 요약을 Slack 채널로 보냅니다.
timeout: PT5M을 걸어 무한 루프를 방지합니다. retry 속성에는 type: exponential, maxAttempts: 4를 명시해 일시적인 순단은 알아서 버티도록 만듭니다.errors 블록을 두고 io.kestra.plugin.notifications.slack.SlackIncomingWebhook 타입을 호출합니다.{{ errorLogs()[0]['message'] }} 구문을 넣어 에러를 유발한 태스크 명세와 원천 스택트레이스를 요약 출력합니다.엔지니어가 로그 대시보드를 찾아 헤매지 않고도 터진 코드 라인을 곧바로 확인할 수 있어서 장애 조치 시간(MTTR)이 짧아집니다.
id: autonomous_error_tracking_pipeline
namespace: production.data.orchestration
description: 일시적 순단을 스스로 방어하고 장애 발생 시 StackTrace 스냅샷을 Slack에 배포하는 워크플로우
tasks:
- id: execute_analytics_ingest
type: io.kestra.plugin.scripts.python.Script
containerImage: python:3.11-slim
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
timeout: PT5M
retry:
type: exponential
interval: PT3S
maxAttempts: 4
maxDuration: PT10M
script: |
import sys
data_integrity_check = False
if not data_integrity_check:
print("[ERROR] CRITICAL: DB Target Schema mismatch detected at row 298. Value 'N/A' violates NULL constraint.", file=sys.stderr)
sys.exit(1)
errors:
- id: dispatch_slack_trace_report
type: io.kestra.plugin.notifications.slack.SlackIncomingWebhook
url: "{{ secret('SLACK_ALERTS_WEBHOOK') }}"
payload: |
{
"channel": "#alerts-platform-ops",
"attachments": [
{
"color": "#E0115F",
"title": "🚨 실시간 데이터 플랫폼 장애 상황 인지 보고서",
"fields": [
{ "title": "네임스페이스 및 플로우 식별 정보", "value": "`{{ flow.namespace }}.{{ flow.id }}`", "short": true },
{ "title": "장애 핵심 코드 및 원천 StackTrace 요약", "value": "`{{ errorLogs()[0]['message'] }}`", "short": false }
]
}
]
}
비정기적인 대규모 벌크 연산이 돌 때 시스템 메모리가 폭증하면서 스케줄러 인스턴스까지 같이 죽는 현상은 흔합니다. Kestra는 워커 그룹(Worker Group) 설정을 통해 시스템 태스크들을 물리 워커 노드의 가용 자원에 따라 맞춤 매핑합니다.
클라우드 워커 연산 전체 비용 은 가용 리소스 할당량 와 작동 시간 의 상호 누적 합에 인프라 고유 단가 비율 를 곱해 계산합니다.
워커 인프라를 띄울 때 명령행 인수를 분리해서 API 호출 전담 소형 풀과 고집적 메모리가 소모되는 전용 high-compute 풀을 나눕니다.
kestra server worker --worker-group core-light 명령으로 구동합니다.kestra server worker --worker-group high-compute 명령으로 따로 띄웁니다.비싼 고사양 클러스터 장비를 24시간 내내 상시 가동하며 낭비하던 인프라 비용을 줄일 수 있습니다. 대기 상태일 때 Celery나 웹서버가 리소스를 상시 점유하는 문제도 해결됩니다.
파이썬 레거시 파이프라인을 걷어내고 Kestra를 안정적으로 안착시키려면 다음 3단계를 순서대로 밟아야 합니다.
운영 중인 파이프라인의 다운타임을 막기 위해 스트랭글러 피그 패턴을 씁니다. Kestra 환경의 Airflow 전용 플러그인(io.kestra.plugin.airflow.dags.TriggerDagRun)을 돌려 새로운 Kestra 스케줄러 콘솔에서 기존 Airflow 상의 데이터 처리를 트리거하고 모니터링 상태를 중앙으로 가져옵니다.
메타데이터 동기화 불일치를 없애기 위해 기존 DAG 파일 내부에 얽혀 있던 비즈니스 가공 파이썬 코드를 Kestra 전용 스크립트 형식의 선언적 YAML 명세로 도려냅니다. 호스트 환경에 Base64 자동 가공 배포 스크립트를 적용해 보안 비밀값들을 변하지 않는 형태로 격리합니다.
kestractl 드라이런 배포 자동화를 적용해 운영 환경의 임의적 웹 조작을 막고 system.readOnly: true 마킹을 강제합니다. worker-group 분리 정책을 펴서 대규모 연산 시 인프라 자원 고갈과 OOM 사태를 끝냅니다. 마지막으로 errorLogs() 기반 자가 진단 스택트레이스 Slack 중계 모니터링 허브를 전역 네임스페이스 계층에 기동합니다.
이 프로세스를 거치면 유지보수 소요 시간이 대폭 줄어듭니다.
결과적으로 단 한 명의 엔지니어 리소스만으로도 수천만 건에 달하는 대용량 데이터 전송 성능을 지탱할 수 있는 인프라가 만들어집니다.