TuBrief
Subscribed Channels
Videos
Community

Airflow에서 Kestra로 넘어가며 겪은 파이프라인 전환 기록

TuBrief Editorial
July 1, 2026
0
Computing/Software

Written with AI assistance from the source video. The video is the authority.

한국어

Related Video

I Tried the Tool Trying to Kill Apache Airflow (Kestra)5:06

I Tried the Tool Trying to Kill Apache Airflow (Kestra)

Better Stack

More from the community

사내 시스템에 llm api 붙일 때 마주하는 현실적인 한계와 대응법

September 13, 2026

레거시 백엔드에 GPT-6 Astra 붙일 때 예산 승인과 보안 통과를 먼저 끝내는 법이 있습니다

September 13, 2026

에이전트끼리 대화하다 6천만 원 청구서가 나오는 이유

September 13, 2026

사내 RAG 벡터 검색에 Okta 권한 필터를 직접 거는 방법

September 13, 2026

브라우저 에이전트에게 내 구글 계정을 통째로 넘기면 안 되는 이유

September 12, 2026

Apple Won the AI Race

September 12, 2026

Comments (0)

Log in to leave a comment

No posts yet

© 2026 . All rights reserved.

TuBrief
Subscribed Channels
Videos
Community
Log in

Airflow에서 Kestra로 넘어가며 겪은 파이프라인 전환 기록

오케스트레이션 정의와 비즈니스 로직이 파이썬 코드로 떡칠된 Apache Airflow 환경은 엔지니어링 리소스가 부족한 스타트업에게 재앙입니다. 관리에 손이 너무 많이 갑니다. 서비스 중단 없이 안전하게 이전을 마치려면 기존 파이프라인을 쪼개어 옮기는 스트랭글러 피그 패턴(Strangler Fig Pattern)을 써야 합니다.

Airflow TaskFlow API 내부 로직을 Kestra의 선언적 YAML 스크립트 태스크로 매핑하는 과정은 생각보다 단순합니다.

  1. Airflow DAG 안에서 데이터를 추출하고 정제하던 파이썬 함수를 독립된 스크립트 블록으로 격리합니다. 실행 환경은 containerImage: python:3.11-slim 기반의 Docker 러너로 지정합니다.
  2. 기존 XCom 직렬화 방식은 버립니다. 가벼운 메타데이터 변수는 Kestra.outputs({"key": value}) JSON 호출로 반환하고, 대용량 데이터프레임은 로컬 디스크에 CSV 파일로 저장합니다. 이 파일을 outputFiles 명세에 등록하면 Kestra 내장 저장소로 들어갑니다.
  3. 다운스트림 태스크의 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 수동 변경을 막는 CI/CD 배포 프로세스

개발자가 웹 콘솔 UI에서 YAML 코드를 직접 고치기 시작하면 형상 관리가 망가집니다. 운영 환경이 언제 터질지 모르는 시한폭탄이 됩니다. 이를 막으려면 운영 환경으로 들어오는 모든 플로우에 system.readOnly: true 라벨을 넣어 원격 UI 에디터를 잠가야 합니다. GitHub Actions와 Kestra CLI 도구인 kestractl을 엮어 배포를 자동화하는 방법입니다.

  1. GitHub 레포지토리에 .github/workflows/deploy.yml 파일을 만들고 on: push와 pull_request 트리거를 겁니다.
  2. 메인 브랜치에 머지하기 전에 kestra-io/github-actions/validate-flows 액션으로 YAML 스키마 구조를 먼저 검증합니다. 원격 API에 영향을 주지 않는 드라이런 방식입니다.
  3. 검증이 끝나면 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와 메모리의 물리 상한선 RtotalR_{\text{total}}Rtotal​이 시스템 총 수용량 R_{\text{worker_capacity}}를 넘지 않도록 제어하는 아키텍처 수식은 다음과 같습니다.

Rtotal=∑i∈tasksCPUi+Memoryi≤Rworker_capacityR_{\text{total}} = \sum_{i \in \text{tasks}} \text{CPU}_{i} + \text{Memory}_{i} \le R_{\text{worker\_capacity}}Rtotal​=i∈tasks∑​CPUi​+Memoryi​≤Rworker_capacity​

네임스페이스별로 concurrencyLimit 정책을 세우고 실행 환경 컨테이너 내부에 물리 리소스 상한을 지정합니다.

  • 경량 수집 영역인 company.data.ingest는 동시 실행 제한을 5로 잡고 cpus: 1.0, memory: 1024Mb 자원을 줍니다.
  • dbt나 대용량 배치를 정제하는 company.data.mart는 자원 독점을 막기 위해 동시 실행 제한을 2로 낮춥니다. 대신 cpus: 4.0, memory: 8192Mb로 연산 범위를 넓혀줍니다.
  • 오픈소스 환경이라면 호스트 수준에서 Base64로 인코딩한 .env_encoded 파일을 도커 컴포즈의 env_file로 연동합니다. 플로우 YAML 안에서는 {{ secret('DATABASE_PROD_PASSWORD') }} 구문만 쓰도록 강제해서 네임스페이스끼리 보안 설정을 훔쳐보지 못하게 막습니다.

데이터가 갑자기 몰려 특정 마트 생성 태스크가 넘치더라도 해당 격리 컨테이너 안에서 OOM(Out Of Memory) Kill 처리됩니다. 제어 엔진 전체가 같이 죽는 불상사를 피할 수 있습니다.


지수 백오프와 StackTrace를 활용한 에러 핸들링

장애 원인도 모른 채 무작정 자동 재시도를 반복하면 스타트업의 지갑이 털립니다. 서드파티 API 쿼터만 고갈되어 2차 장애로 이어지기 십상입니다. 일시적인 네트워크 결함과 데이터 자체의 스키마 에러를 분리해서 대응해야 합니다.

예외 유형 장애 원인 예시 대응 액션 및 복구 전략 알림 제어 방식
일시적 네트워크 예외 외부 API 순단 (HTTP 503), DB 커넥션 병목 retry 블록 기반 지수 백오프 가동 (PT2S→PT4S→PT8SPT2S \rightarrow PT4S \rightarrow PT8SPT2S→PT4S→PT8S) 최대 임계치까지 재시도 실패 시에만 장애 채널 알림
구조적 데이터 결함 정합성 불일치 (Null Value 유입), 스키마 위반 allowFailed: true 우회 처리 또는 Fail 즉시 전파 StackTrace 분석 후 실시간 On-call 핫라인 전송
비핵심 태스크 장애 마케팅 메일 발송 지연, 단순 통계 누락 allowWarning: true 지정으로 후속 작업 진행 경고 알림 전송 후 워크플로우는 정상 완료 처리

태스크 실패가 터지면 전역 에러 제어기를 돌려 Kestra 내장 함수 errorLogs()로 추출한 StackTrace 요약을 Slack 채널로 보냅니다.

  1. 개별 태스크 블록에 timeout: PT5M을 걸어 무한 루프를 방지합니다. retry 속성에는 type: exponential, maxAttempts: 4를 명시해 일시적인 순단은 알아서 버티도록 만듭니다.
  2. 플로우 최하단에 errors 블록을 두고 io.kestra.plugin.notifications.slack.SlackIncomingWebhook 타입을 호출합니다.
  3. Slack 페이로드 message 필드에 {{ 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) 설정을 통해 시스템 태스크들을 물리 워커 노드의 가용 자원에 따라 맞춤 매핑합니다.

클라우드 워커 연산 전체 비용 CtotalC_{\text{total}}Ctotal​은 가용 리소스 할당량 RiR_iRi​와 작동 시간 TiT_iTi​의 상호 누적 합에 인프라 고유 단가 비율 WrateW_{\text{rate}}Wrate​를 곱해 계산합니다.

Ctotal=∑i=1N(Ri×Ti)×WrateC_{\text{total}} = \sum_{i=1}^{N} (R_{i} \times T_{i}) \times W_{\text{rate}}Ctotal​=i=1∑N​(Ri​×Ti​)×Wrate​

워커 인프라를 띄울 때 명령행 인수를 분리해서 API 호출 전담 소형 풀과 고집적 메모리가 소모되는 전용 high-compute 풀을 나눕니다.

  1. 가벼운 제어를 담당하는 일반 인스턴스는 kestra server worker --worker-group core-light 명령으로 구동합니다.
  2. 고사양 메모리가 필요한 파티션 처리 전용 서버는 kestra server worker --worker-group high-compute 명령으로 따로 띄웁니다.
  3. Kestra 백엔드의 JDBC 기반 분산 구조 덕분에 각 워커 그룹 인스턴스들이 중앙 트랜잭션 큐 저장 테이블을 실시간으로 폴링합니다. FIFO 기반으로 배분되므로 유휴 주기가 짧고 연산 속도가 빠른 고사양 워커 노드가 테이블 조회를 끝내고 자율적으로 큰 부하를 가져갑니다.

비싼 고사양 클러스터 장비를 24시간 내내 상시 가동하며 낭비하던 인프라 비용을 줄일 수 있습니다. 대기 상태일 때 Celery나 웹서버가 리소스를 상시 점유하는 문제도 해결됩니다.


스타트업을 위한 3단계 마이그레이션 행동 지침

파이썬 레거시 파이프라인을 걷어내고 Kestra를 안정적으로 안착시키려면 다음 3단계를 순서대로 밟아야 합니다.

1단계: 하이브리드 제어권 확보 (1~2주차)

운영 중인 파이프라인의 다운타임을 막기 위해 스트랭글러 피그 패턴을 씁니다. Kestra 환경의 Airflow 전용 플러그인(io.kestra.plugin.airflow.dags.TriggerDagRun)을 돌려 새로운 Kestra 스케줄러 콘솔에서 기존 Airflow 상의 데이터 처리를 트리거하고 모니터링 상태를 중앙으로 가져옵니다.

2단계: 비즈니스 격리 리팩토링 및 비밀값 정비 (3~4주차)

메타데이터 동기화 불일치를 없애기 위해 기존 DAG 파일 내부에 얽혀 있던 비즈니스 가공 파이썬 코드를 Kestra 전용 스크립트 형식의 선언적 YAML 명세로 도려냅니다. 호스트 환경에 Base64 자동 가공 배포 스크립트를 적용해 보안 비밀값들을 변하지 않는 형태로 격리합니다.

3단계: 리소스 분산 및 온콜 경보 자동화 (5주차 이후)

kestractl 드라이런 배포 자동화를 적용해 운영 환경의 임의적 웹 조작을 막고 system.readOnly: true 마킹을 강제합니다. worker-group 분리 정책을 펴서 대규모 연산 시 인프라 자원 고갈과 OOM 사태를 끝냅니다. 마지막으로 errorLogs() 기반 자가 진단 스택트레이스 Slack 중계 모니터링 허브를 전역 네임스페이스 계층에 기동합니다.

이 프로세스를 거치면 유지보수 소요 시간이 대폭 줄어듭니다.

Tmaint_new=Tmaint_legacy×0.6T_{\text{maint\_new}} = T_{\text{maint\_legacy}} \times 0.6Tmaint_new​=Tmaint_legacy​×0.6

결과적으로 단 한 명의 엔지니어 리소스만으로도 수천만 건에 달하는 대용량 데이터 전송 성능을 지탱할 수 있는 인프라가 만들어집니다.