call_api (SimpleHttpOperator)
--> process_api_data (PythonOperator)
--> analyze_data (PythonOperator)
--> send_email (EmailOperator)
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
)
사전에 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
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 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로 전송
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 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 인증때문인것 같다..