Summary
- Connection은 Airflow UI 화면에서 등록한 커넥션 정보다.
- Hook은 Airflow 외부 솔루션 기능을 사용할 수 있도록 미리 구현된 메서드를 가진 클래스로, Connection 정보를 통해 생성되어 접속정보가 코드상 노출되지 않는다.
- Hook은 task를 만들지 못하므로 커스텀 오퍼레이터 또는 python 오퍼레이터 내 함수로 사용한다.
Connection
Connection은 Airflow UI 화면에서 등록한 커넥션 정보다.
| key | value |
|---|---|
| Connection id | conn=db-postgres-custom |
| Connection Type | Postgres |
| Host | 172.28.0.3 |
| Database | haejun |
| Login | haejun |
| Password | **** |
| Port | 5432 |
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.connbulk_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 여부 선택, 특수문자 제거 로직을 추가할 수 있다.