Summary

  • Connection은 Airflow UI 화면에서 등록한 커넥션 정보다.
  • Hook은 Airflow 외부 솔루션 기능을 사용할 수 있도록 미리 구현된 메서드를 가진 클래스로, Connection 정보를 통해 생성되어 접속정보가 코드상 노출되지 않는다.
  • Hook은 task를 만들지 못하므로 커스텀 오퍼레이터 또는 python 오퍼레이터 내 함수로 사용한다.

Connection

Connection은 Airflow UI 화면에서 등록한 커넥션 정보다.

keyvalue
Connection idconn=db-postgres-custom
Connection TypePostgres
Host172.28.0.3
Databasehaejun
Loginhaejun
Password****
Port5432

Hook

Hook은 Airflow 외부 솔루션 기능을 사용할 수 있도록 미리 구현된 메서드를 가진 클래스다. Connection 정보를 통해 생성되므로 접속정보가 코드상 노출되지 않는다. task를 만들지 못하므로, 커스텀 오퍼레이터 또는 python 오퍼레이터 내 함수로 사용한다.

get_conn()

Connection 정보를 가져와 DB에 연결한다.

def get_conn(self) -> connection:
    conn_id = getattr(self, self.conn_name_attr)
    conn = deepcopy(self.connection or self.get_connection(conn_id))
    conn_args = {
        "host": conn.host,
        "user": conn.login,
        "password": conn.password,
        "dbname": self.database or conn.schema,
        "port": conn.port,
    }
    self.conn = psycopg2.connect(**conn_args)
    return self.conn

bulk_load

def bulk_load(self, table: str, tmp_file: str) -> None:
    self.copy_expert(f"COPY {table} FROM STDIN", tmp_file)

기존 bulk_load 단점

  • 구분자가 \t으로 고정된다.
  • Header까지 포함되어 업로드된다.
  • 특수문자로 인해 파싱이 안될 경우 에러가 발생한다.

Custom Hook

기존 Hook의 단점을 개선하려면 BaseHook 클래스를 상속받아 Custom Hook을 작성한다. 구분자 유형 입력, Header 여부 선택, 특수문자 제거 로직을 추가할 수 있다.