Airflow Api호출 후 이메일 전송

BAO.DE·2025년 9월 25일

Apache Airflow

목록 보기
18/20

실행흐름

call_api (SimpleHttpOperator)
--> process_api_data (PythonOperator)
--> analyze_data (PythonOperator)
--> send_email (EmailOperator)

HttpOperator

from airflow.providers.http.operators.http import HttpOperator

Airlfow에서 api 호출할때 사용하는 operator이다

사용방법

Connections에 http connections을 사전 등록후 해당 id를
http_conn_id = 'http_default' 형식으로 지정

이렇게 하면 DAG 코드에서는 민감 정보나 반복되는 URL을 하드코딩하지 않아도 됨.


api_task = SimpleHttpOperator(
    task_id='call_api',
    http_conn_id='http_default',  # Airflow UI → Admin → Connections에서 등록 필요
    endpoint='posts',             # 호출할 API 경로
    method='GET',
    headers={"Content-Type": "application/json"},
    response_filter=lambda response: response.text,  # XCom으로 문자열 반환
    dag=dag
)

ping 테스트

사전에 dag을 올리기전에 통신 잘되는지 test

import requests
import json

url = "https://jsonplaceholder.typicode.com/posts"
try:
    response = requests.get(url, timeout=5)
    response.raise_for_status()  # 200 OK 아니면 예외 발생
    data = response.json()
    print(f" API 호출 성공, 글 개수: {len(data)}")
except requests.exceptions.RequestException as e:
    print(f" API 호출 실패: {e}")

✅ API 호출 성공, 글 개수: 100

ping 테스트 ( Airflow )

local에서 ping 테스트 확인후 Airflow webserver에서 ping테스트

def print_response(response):
    print("✅ HTTP Status:", response.status_code)
    print("✅ Response Body (first 300 chars):", response.text[:300])  # 일부만 출력
    return response.text



ping_api = HttpOperator(
    task_id='ping_api',
    http_conn_id='http_default',  
    endpoint='posts',             
    method='GET',
    response_check=lambda response: response.status_code == 200,  # 상태 200 체크
    response_filter=print_response,  # 호출 후 출력
    dag=dag
)

connection에 endpoint까지 집어넣어서 애먹었다 ㅠ URL만 넣도록하자.
Http Status : 200 응답 확인

Json contents Email로 전송하기

랜덤하게 선택된 json id 1개 EmailOperator발송하기

ping_api >> process_task

request api -> response -> Xcom.pull 진행

api 응답 확인

{
"userId": 3,
"id": 27,
"title": "quasi id et eos tenetur aut quo autem",
"body": "eum sed dolores ipsam sint possimus debitis occaecati\ndebitis qui qui et\nut placeat enim earum aut odit facilis\nconsequatur suscipit necessitatibus rerum sed inventore temporibus consequatur"
},

여기서 title과 body를 받아와서 email로 전송

선택된 id 글 요약 DAG

ping_api >> process_task >> email_prepare_task

def prepare_email(**context):
    ti = context['task_instace'] ## 현재객체 key
    selected_post = ti.xcom_pull(task_ids='process_api_data_task',key='selected_post')
    email_body = f"""
    📌 랜덤 선택된 글 요약:
    ID: {selected_post['id']}
    Title: {selected_post['title']}
    Body: {selected_post['body']}
    """
    ti.xcom_push(key='email_body', value=email_body)
    print(email_body)

task instance 받아와서 앞에 task에서 xcom_pull

selected_post에 객체를 담아서 파싱

EmailOperator

마지막으로 EmailOperator task 달아서 마무리
airflow.cfg [smtp] 설정필요

Airflow 구성 Gmail

AIRFLOWSMTPSMTP_HOST smtp.gmail.com
AIRFLOWSMTPSMTP_STARTTLS True
AIRFLOWSMTPSMTP_SSL Fals
AIRFLOWSMTPSMTP_USER your_email@gmail.com
AIRFLOWSMTPSMTP_PASSWORD password
AIRFLOWSMTPSMTP_PORT 587
AIRFLOWSMTPSMTP_MAIL_FROM your_email@gmail.com

send_email_task = EmailOperator(
task_id='send_email',
to='your_email@example.com',
subject='랜덤 글 요약',
html_content="{{ ti.xcom_pull(task_ids='prepare_email_task', key='email_body') }}",
dag=dag
)

요렇게 작성해주고 마찬가지로 email body를 key로 설정

email 발송은 실패 ㅠ SSL 인증때문인것 같다..

0개의 댓글