apache-airflow-providers-clickhousedb로 마이그레이션하는 방법을 설명합니다. provider 자체의 동작 방식은
Apache Airflow를 ClickHouse에 연결하기를 참조하십시오.
왜 네이티브 provider인가?
provider는 Airflow가 서드파티 시스템과 통합되는 방식입니다. provider는 Airflow 생태계의 나머지 구성 요소와 함께 릴리스되고 테스트되며 문서화되며, 이 provider는 Airflow 커뮤니티가 ClickHouse 팀과 함께 유지 관리합니다. 또한 커뮤니티가 유지 관리하는 driver가 아니라, ClickHouse가 직접 개발하고 지원하는 Python client인 clickhouse-connect를 기반으로 구축되었습니다. 따라서 새로운 서버 기능과 수정 사항이 공식 지원 경로를 통해 Airflow 사용자에게 전달됩니다. provider로 전환하면 공식적으로 관리되는 package와 표준common.sql operator 및 sensor,
그리고 다른 모든 데이터베이스와 마찬가지로 Airflow UI에 표시되는 connection 유형을 사용할 수 있습니다.
두 package의 차이는 import 경로에 그치지 않습니다. plugin은 clickhouse-driver를 사용해
네이티브 TCP 프로토콜로 ClickHouse와 통신합니다. provider는 clickhouse-connect를 사용해
**HTTP(S)**로 통신하며, ClickHouse 전용 operator를 별도로 제공하는 대신 범용 common.sql operator에
연결됩니다.
무엇이든 변경하기 전에 가이드 전체를 한 번 읽어 보십시오. 특히 connection 변경은 모든 DAG에 동시에
영향을 미칩니다.
한눈에 보기
1단계: 사전 요구 사항 확인 및 설치
이 provider를 사용하려면 Airflow 2.11 이상과apache-airflow-providers-common-sql 1.32.0 이상이 필요합니다. 그보다 이전 버전을 사용하고 있다면 먼저 Airflow를 업그레이드하십시오.
2단계: 연결 업데이트
이 단계를 건너뛰면 문제가 발생합니다. 기존 ClickHouse 연결은 모두 네이티브 포트를 가리키고 있지만, provider는 HTTP 포트가 필요합니다.
ClickHouse가 firewall 또는 load balancer 뒤에 있다면, 전환하기 전에 worker에서 HTTP 포트에
연결할 수 있는지 확인하십시오. 또한 서버에서 HTTP 인터페이스가 활성화되어 있는지(서버 구성의
http_port 또는
https_port) 확인하십시오. ClickHouse Cloud는 8443 포트에서만
HTTPS를 제공합니다.
URI 형태로 저장된 연결(clickhouse://user:pass@host:9000/db?secure=true)은 Airflow가 URI scheme에서
연결 유형을 가져오므로 이미 clickhouse 연결 유형으로 설정되어 있습니다. 이 경우 포트와 extra keys만 변경하면 됩니다.
쿼리 문자열 값은 JSON으로 파싱되므로 secure=true는 Boolean 값으로 유지됩니다.
연결 extra 옵션
플러그인은extra의 모든 키를 clickhouse_driver.Client에 그대로 전달했기 때문에, 연결에는 어떤 clickhouse-driver keyword argument든 포함될 수 있었습니다. provider는 고정된 키 집합만 읽고, 나머지는 client_kwargs를 통해 전달합니다. 다음과 같이 변환하십시오:
변경 전:
3단계: import 교체
ClickHouse 접두사가 붙은 common.sql 래퍼는 플러그인의 후크를 주입하기 위한 용도로만 존재했습니다. provider가 clickhouse connection 유형을 등록하기 때문에, 접두사가 없는 common.sql 클래스는 connection에서 후크를 직접 확인합니다. 이러한 래퍼를 사용했다면, 마이그레이션은 대개 import 줄을 수정하고 ClickHouse 접두사를 제거한 뒤 conn_id를 명시적으로 전달하는 것으로 끝납니다(7단계 참조).
4단계: ClickHouseOperator에서 SQLExecuteQueryOperator로
이전:
다중 SQL 문 결과
plugin은 마지막 SQL 문의 결과를 XCom에 푸시했습니다.SQLExecuteQueryOperator는 sql이 리스트일 때
SQL 문마다 하나씩 결과를 반환하므로, plugin이 [(12345.0,)]를 푸시했던 것과 달리 위 예시는
[[], [(12345.0,)]]를 푸시합니다. 다음 중 하나를 선택하십시오.
- 다운스트림의
xcom_pull이 마지막 원소를 가져오도록 변경합니다. - SQL 문들을
;로 구분한 단일 문자열로 전달하고split_statements=True로 설정합니다. 기본값인return_last=True와 함께 사용하면 연산자가 마지막 SQL 문의 행만 푸시하므로 plugin과 동일하게 동작합니다.
parameters를 통해 행 목록을 삽입하던 ClickHouseOperator는 SQLInsertRowsOperator로 대체됩니다:
columns를 전달하십시오. 전달하지 않으면 연산자가 SQLAlchemy를 통해 테이블을 조회합니다.
컬럼 타입 유지하기
with_column_types=True는 (rows, [(name, type), ...])를 반환했습니다. handler를 사용하면 이를 재현할 수 있으며, clickhouse-connect cursor는 cursor.description에 ClickHouse 타입 이름을 제공합니다:
Step 5: ClickHouseHook.execute에서 DbApiHook 메서드로
플러그인의 후크는 clickhouse_driver.Client.execute를 그대로 반영한 execute 메서드 하나만 노출했습니다.
provider의 후크는 DbApiHook이므로, 다른 모든 SQL provider와 동일한 표준 메서드를 사용할 수 있습니다:
run, get_records, get_first, get_pandas_df, get_df, insert_rows, test_connection.
clickhouse_conn_id 및 database 생성자 인수는 변경되지 않았습니다.
fetch_all_handler 및 그 외 handler는
airflow.providers.common.sql.hooks.handlers에서 가져옵니다.
플러그인 시절 코드에서 가장 흔한 후크 사용 패턴은 bulk insert입니다. DB-API 커서가 행을 SQL 문자열에
포맷해 넣으려 하기 때문에, 단순한 run 호출로는 처리할 수 없습니다. 대신 native insert를
사용하십시오:
이전:
bulk_insert_rows에는 column_names가 필요합니다. batch_size는 선택 사항이며, 매우 큰 입력에서 메모리 사용량을 제한하는 역할을 합니다. 범용 메서드인 insert_rows(table, rows, target_fields=[...], executemany=True)도 최종적으로 네이티브 삽입으로 처리되지만, executemany=True일 때만 그렇고 기본값에서는 행마다 HTTP request를 한 번씩 전송합니다.
DB-API 계층에서 다루지 못하는 작업은 get_client()를 사용하면 되며, 이 메서드는 Airflow connection 정보로 구성된 raw clickhouse-connect 클라이언트를 반환합니다. plugin이 노출했던 모든 clickhouse-driver 전용 인수를 대체하는 수단입니다:
Step 6: ClickHouseSensor에서 SqlSensor로
호출 가능 객체(callable)의 입력이 바뀌는 유일한 대체 작업입니다.
plugin이 전체 결과 집합을 그대로 넘겨주었기 때문에, 기존에 동작하던 센서 코드는 결과에 인덱스로 접근합니다. 마이그레이션할 때는 이러한 인덱싱을 제거하십시오.
변경 전:
selector=lambda row: row를 전달하십시오. 모든 행이 필요하다면, 쿼리가 단일 불리언 값이나 개수를 반환하도록 검사 로직을 SQL로 작성하십시오.
7단계: common.sql 래퍼 계열
ClickHouseSQLExecuteQueryOperator, ClickHouseSqlSensor 및 그 외 ClickHouse 접두사가 붙은 래퍼를 사용하던 코드는 손볼 부분이 가장 적습니다:
- import를
common.sql모듈로 변경하고 클래스 이름에서ClickHouse접두사를 제거하십시오. conn_id를 명시적으로 전달하십시오. 래퍼는conn_id가 없거나None이면clickhouse_default로 처리했지만,common.sql클래스에는 기본값이 없어 지정하지 않으면 실패합니다.default_args={"conn_id": "clickhouse_default"}를 사용하면 DAG 전체에 적용됩니다.- 연산자의
database=와 센서의hook_params={"schema": ...}는 계속 동작합니다. provider의 후크는schema를database의 별칭으로 취급합니다. ClickHouseDbApiHook은ClickHouseHook으로 바뀝니다. 생성자 인수schema는 여전히 별칭으로 허용되며, ClickHouse의 네이티브 표기는database이므로 둘 다 지정된 경우database가 우선합니다.- connection에는 2단계에서 설명한 포트 및 extras 변경이 여전히 필요합니다. 래퍼 역시 native protocol을 사용했기 때문입니다.
검토해야 할 동작 차이
코드가 정상적으로 컴파일된 뒤에도 런타임 동작이 몇 가지 달라집니다. INSERT 작업의 XCom 값. 플러그인은clickhouse-driver가 반환한 값을 그대로 푸시했으며, 파라미터를 사용한 VALUES 삽입에서는 이 값이 삽입된 행 수였습니다. 반면 provider는 행을 반환하지 않는 SQL 문에 대해 빈 결과 집합을 푸시합니다. 따라서 XCom에서 행 수를 읽던 다운스트림 작업은 후속 SELECT count() 등 다른 방법으로 행 수를 가져와야 합니다.
세션 설정과 SET 문. 두 패키지 모두 단일 연결에서 여러 SQL 문 목록을 실행하며, clickhouse-connect는 기본적으로 클라이언트당 세션을 생성하므로 목록 앞부분의 SET 문은 이후 SQL 문에도 계속 적용됩니다. 그래도 session_settings를 사용하는 것이 좋습니다. 명시적이고 템플릿을 적용할 수 있으며, 서버나 중간 프록시가 세션을 유지하는지 여부와 무관하게 동일하게 동작합니다. DAG가 SET에 의존한다면 해당 환경에서 동작을 직접 확인하십시오.
예외. 오류는 이제 clickhouse_driver.errors.ServerException 및 NetworkError가 아니라 clickhouse_connect.driver.exceptions.DatabaseError, OperationalError 또는 ProgrammingError로 발생합니다. 예외 타입을 검사하는 except 절, on_failure_callback 코드, 재시도 로직을 업데이트하십시오.
압축. clickhouse-driver는 compression을 설정하지 않으면 압축을 사용하지 않았지만, clickhouse-connect는 기본적으로 HTTP 응답 압축을 활성화하고 서버와 알고리즘을 협상합니다. 이전 동작으로 되돌리려면 연결 extra에 "compress": false를 설정하십시오. 네이티브 압축을 위해 플러그인에 필요했던 clickhouse-cityhash 패키지는 더 이상 필요하지 않으며, lz4는 clickhouse-connect 의존성으로 계속 설치됩니다.
타입 매핑. 두 driver 모두 네이티브 Python 타입을 반환하지만, 서로 코드 베이스가 다릅니다. 시간대가 지정된 DateTime64, Decimal, UUID, Nullable 컬럼, 중첩된 Array 또는 Map 값의 정확한 타입에 의존하는 작업을 검토하십시오. 특히 결과가 XCom에 푸시되어 다운스트림에서 사용되는 경우 주의가 필요합니다.
system.query_log에서의 쿼리 식별. 쿼리는 이제 HTTP 인터페이스를 통해 전달되므로 interface = 1이 아닌 interface = 2로 표시되며, system.query_log의 http_user_agent 컬럼에는 Airflow 및 provider 버전과, 설정된 경우 client_name extra가 함께 담깁니다. native protocol이나 clickhouse-driver 클라이언트 이름을 기준으로 필터링하던 모니터링은 업데이트해야 합니다. 행을 반환하지 않는 SELECT는 항목을 하나 더 생성합니다. DB-API 커서가 컬럼 메타데이터를 얻기 위해 SELECT * FROM (...) LIMIT 0을 실행하기 때문입니다.
HTTP를 통한 타임아웃. send_receive_timeout은 이제 HTTP 읽기 타임아웃이며, worker와 ClickHouse 사이에 있는 프록시나 로드 밸런서는 요청에 각자의 idle timeout을 적용합니다. native protocol에서 수 분 동안 실행되던 SQL 문은 이러한 제한값을 높여야 할 수 있습니다.
연결 처리. 후크는 run 또는 get_client 호출마다 clickhouse-connect 클라이언트를 생성하며, 이는 플러그인이 execute마다 새 native connection을 열었던 방식과 같습니다. 클라이언트들은 프로세스 전역 HTTP 연결 풀을 공유하므로, get_client()로 얻은 클라이언트에 close()를 호출하는 것은 필수는 아니지만 권장되는 습관입니다. 풀은 작업 프로세스가 종료될 때 해제됩니다. 클라이언트는 context manager이므로 with hook.get_client() as client: 형태가 가장 깔끔합니다.
체크리스트
- Airflow 버전이 2.11 이상입니다.
- worker에서 HTTP 포트에 연결할 수 있으며, HTTP 엔드포인트용 TLS 인증서가 유효합니다.
- 모든 ClickHouse connection: 타입은
clickhouse, 포트는8123또는8443, extras는 2단계에 따라 변환했으며,airflow connections test로 검증했습니다. - 3단계에 따라 import를 대체했으며,
common.sql래퍼에서ClickHouse프리픽스를 제거했습니다. - 연산자와 센서에서
clickhouse_conn_id를conn_id로 이름을 변경했고, 플러그인의 기본값에 의존하던 모든 작업에conn_id를 설정했습니다. settings=를hook_params={"session_settings": ...}안으로 옮겼습니다.hook.execute("INSERT ... VALUES", rows)를bulk_insert_rows로 대체했고,ClickHouseOperator(parameters=rows)를SQLInsertRowsOperator로 대체했습니다.- 센서 호출 가능 객체을
result[0][0]대신 cell 값 자체를 사용하도록 수정했습니다. with_column_types,external_tables,columnar,query_id,types_check를 사용하던 부분을handler또는get_client()로 재작성했습니다.- 다중 SQL 문 및 INSERT 작업의 XCom을 사용하는 다운스트림 소비자를 검토했습니다.
clickhouse_driver예외를 처리하던 코드를 업데이트했습니다.airflow-clickhouse-plugin과clickhouse-driver를 제거했습니다.