전체 그래프
Airflow

Connections & Hooks

data-engineeringtool/airflowconnectionhook

상위: Airflow

요약

Connection은 Airflow가 외부 시스템(DB, API, S3 등)과 연결하기 위한 설정 정보입니다. Hook은 Connection을 사용하여 실제 데이터 작업을 수행하는 인터페이스 클래스입니다. Web UI, CLI, 환경변수로 Connection을 관리하며, Hook을 통해 Operator 내부에서 외부 시스템과 상호작용합니다.

Connection이란?

  • Airflow가 **외부 시스템(DB, API, S3 등)**과 연결하기 위한 정보
  • 호스트, 포트, 사용자명, 비밀번호 등 저장
  • Web UI, CLI, 환경변수로 설정 가능
  • 메타 DB에 암호화되어 저장

Hook이란?

  • Connection을 사용해 실제 데이터 작업을 수행하는 인터페이스 클래스
  • Operator 내부에서 사용
  • 외부 시스템과 상호작용
  • 재사용 가능한 코드 제공

Connection & Hook 관계

연결 대상Connection Type사용 Hook
MySQLMySQLMySqlHook
PostgreSQLPostgresPostgresHook
REST APIHTTPHttpHook
AWS S3AWSS3Hook
Google CloudGoogle CloudGCSHook
SSHSSHSSHHook

Connection 설정

Web UI에서 설정

  1. Admin → Connections → Add Connection

Connections & Hooks

  1. Connection IDConnection Type 설정

Connections & Hooks

주요 필드:

  • Connection Id: 고유 식별자
  • Connection Type: 연결 유형 (Postgres, MySQL, HTTP 등)
  • Host: 호스트 주소
  • Schema: 데이터베이스 이름
  • Login: 사용자명
  • Password: 비밀번호
  • Port: 포트 번호

CLI로 설정

# Connection 추가
airflow connections add my_postgres_conn \
    --conn-type postgres \
    --conn-host localhost \
    --conn-login airflow \
    --conn-password airflow \
    --conn-port 5432 \
    --conn-schema airflow_db

# Connection 목록 조회
airflow connections list

# Connection 삭제
airflow connections delete my_postgres_conn

환경변수로 설정

export AIRFLOW_CONN_MY_POSTGRES='postgres://airflow:airflow@localhost:5432/airflow_db'

URI 형식: {conn_type}://{login}:{password}@{host}:{port}/{schema}

Hook 사용

PostgresHook 예제

from airflow.providers.postgres.hooks.postgres import PostgresHook

def query_database():
    # Connection 사용
    pg_hook = PostgresHook(postgres_conn_id="my_postgres_conn")
    
    # SQL 실행
    records = pg_hook.get_records("SELECT * FROM users LIMIT 10")
    
    for record in records:
        print(record)

task = PythonOperator(
    task_id='query_db',
    python_callable=query_database
)

MySqlHook 예제

from airflow.providers.mysql.hooks.mysql import MySqlHook

def use_mysql():
    mysql_hook = MySqlHook(mysql_conn_id="my_mysql_conn")
    
    # DataFrame으로 조회
    df = mysql_hook.get_pandas_df(sql="SELECT * FROM users")
    print(df.head())

HttpHook 예제

from airflow.providers.http.hooks.http import HttpHook

def call_api():
    http_hook = HttpHook(method='GET', http_conn_id='my_api_conn')
    
    # API 호출
    response = http_hook.run('/api/users')
    data = response.json()
    
    print(data)

S3Hook 예제

from airflow.providers.amazon.aws.hooks.s3 import S3Hook

def upload_to_s3():
    s3_hook = S3Hook(aws_conn_id='aws_default')
    
    # 파일 업로드
    s3_hook.load_file(
        filename='/tmp/data.csv',
        key='data/data.csv',
        bucket_name='my-bucket'
    )

Hook 주요 메서드

PostgresHook / MySqlHook

hook = PostgresHook(postgres_conn_id='conn_id')

# 레코드 조회
records = hook.get_records("SELECT * FROM table")

# 단일  조회
value = hook.get_first("SELECT COUNT(*) FROM table")

# Pandas DataFrame
df = hook.get_pandas_df("SELECT * FROM table")

# SQL 실행
hook.run("INSERT INTO table VALUES (1, 'name')")

# 대량 삽입
hook.insert_rows(table='users', rows=[(1, 'Alice'), (2, 'Bob')])

HttpHook

hook = HttpHook(method='GET', http_conn_id='api_conn')

# API 호출
response = hook.run(endpoint='/api/resource')

# POST 요청
response = hook.run(
    endpoint='/api/create',
    data={'name': 'Alice'},
    headers={'Content-Type': 'application/json'}
)

S3Hook

hook = S3Hook(aws_conn_id='aws_default')

# 파일 업로드
hook.load_file(filename='local.csv', key='s3/path.csv', bucket_name='bucket')

# 파일 다운로드
hook.download_file(key='s3/path.csv', bucket_name='bucket', local_path='/tmp')

# 파일 존재 확인
exists = hook.check_for_key(key='s3/path.csv', bucket_name='bucket')

# 파일 삭제
hook.delete_objects(bucket='bucket', keys=['file1.csv', 'file2.csv'])

Operator와 Hook

Operator는 내부적으로 Hook 사용

# PostgresOperator 내부
class PostgresOperator(BaseOperator):
    def execute(self, context):
        hook = PostgresHook(postgres_conn_id=self.postgres_conn_id)
        hook.run(self.sql)

직접 Hook 사용 vs Operator 사용

#  Operator 사용 (간단한 경우)
postgres_task = PostgresOperator(
    task_id='run_query',
    postgres_conn_id='my_conn',
    sql='SELECT * FROM users'
)

#  Hook 사용 (복잡한 로직)
def custom_logic():
    hook = PostgresHook(postgres_conn_id='my_conn')
    records = hook.get_records('SELECT * FROM users')
    
    # 복잡한 처리 로직
    for record in records:
        # 데이터 가공, 조건 분기 
        ...

task = PythonOperator(
    task_id='custom_task',
    python_callable=custom_logic
)

보안 Best Practices

Connection에 민감 정보 저장

#  Connection 사용 (암호화됨)
hook = PostgresHook(postgres_conn_id='secure_conn')

#  Variable에 비밀번호 저장 (보안 취약)
password = Variable.get('db_password')

Fernet Key 설정

# Fernet key 생성
python -c "from cryptography.fernet import Fernet; print(Fernet.generate_key().decode())"

# airflow.cfg 또는 환경변수
export AIRFLOW__CORE__FERNET_KEY='your-fernet-key'

주의사항

Connection ID 명명 규칙

  • 명확하고 설명적인 이름 사용
  • {service}_{environment}_{purpose} 형식 권장
  • 예: postgres_prod_users, s3_dev_data

Connection 테스트

# Connection 테스트
from airflow.hooks.base import BaseHook

try:
    conn = BaseHook.get_connection('my_conn')
    print(f"Host: {conn.host}")
    print(f"Login: {conn.login}")
except:
    print("Connection not found!")

다음 단계

  • Variable - 전역 변수 관리
  • Operators - Operator 종류와 사용법
  • Spark Integration - Spark Connection 설정

참고 자료