apache-airflow-providers-clickhousedb へ移行する手順を説明します。プロバイダ自体の動作については、Apache Airflow と ClickHouse の接続を参照してください。
なぜネイティブプロバイダなのか?
プロバイダは、Airflow がサードパーティ製システムと連携するための仕組みです。Airflow エコシステムの他の構成要素と同じタイミングでリリース・テスト・ドキュメント化されており、本プロバイダは Airflow コミュニティと ClickHouse チームが共同でメンテナンスしています。コミュニティがメンテナンスするドライバーではなく、ClickHouse 自身が開発・サポートする Python クライアントである clickhouse-connect を基盤としています。そのため、新しいサーバー機能や修正は、サポートされた経路を通じて Airflow ユーザーのもとに届きます。プロバイダに移行すれば、公式の提供元を持つパッケージ、標準のcommon.sql オペレーターとセンサー、そして他のデータベースと同様に Airflow UI に表示される接続タイプが利用できます。
2 つのパッケージの違いは import パスだけではありません。プラグインは clickhouse-driver を用いて ネイティブ TCP プロトコル で ClickHouse と通信します。一方プロバイダは clickhouse-connect を用いて HTTP(S) で通信し、ClickHouse 固有のオペレーターを同梱する代わりに、汎用の common.sql オペレーターに組み込まれます。
変更を加える前に、このガイドを一度通読してください。特に接続の変更は、すべての DAG に同時に影響します。
概要
ステップ 1: 前提条件の確認とインストール
このプロバイダには Airflow 2.11 以降とapache-airflow-providers-common-sql 1.32.0 以降が必要です。これより古いバージョンを使用している場合は、まず Airflow をアップグレードしてください。
ステップ 2: connection を更新する
このステップを省略すると動作しなくなります。既存の ClickHouse connection はすべてネイティブポートを指しているため、provider が必要とする HTTP ポートに変更する必要があります。
ClickHouse がファイアウォールや load balancer の背後にある場合は、切り替える前に worker から HTTP ポートに到達できることを確認してください。あわせて、server 側で HTTP インターフェイスが有効になっているか (server configuration の
http_port または https_port) も確認してください。ClickHouse Cloud では HTTPS は 8443 のみで公開されています。
URI として保存されている connection (clickhouse://user:pass@host:9000/db?secure=true) の場合、Airflow が URI のスキームから判別するため、接続タイプはすでに clickhouse になっています。変更が必要なのはポートと extra のオプションだけです。クエリ文字列の値は JSON としてパースされるため、secure=true はブール値のまま扱われます。
接続の extra
プラグインはextra 内のすべてのキーをそのまま clickhouse_driver.Client に渡していたため、接続には clickhouse-driver の任意のキーワード引数を指定できました。プロバイダは決められたキーのみを読み取り、それ以外は client_kwargs 経由で転送します。次のように読み替えてください。
変更前:
ステップ3: import の置き換え
ClickHouse プレフィックスが付いた common.sql の wrapper は、プラグインの hook を注入するためだけに存在していたものです。プロバイダが clickhouse 接続タイプを登録するため、プレフィックスなしの common.sql のクラスは connection から hook を自力で解決できます。これらの wrapper を使用していた場合、移行作業は通常、import 行の書き換え、ClickHouse プレフィックスの削除、conn_id の明示的な指定だけで済みます (ステップ7を参照) 。
ステップ4: ClickHouseOperator から SQLExecuteQueryOperator へ
変更前:
複数ステートメントの結果
このpluginは最後のステートメントの結果をXComにプッシュしていました。SQLExecuteQueryOperatorはsqlがリストの場合、ステートメントごとに1つの結果を返します。そのため、上記の例では[[], [(12345.0,)]]がプッシュされますが、pluginでは[(12345.0,)]がプッシュされていました。次のいずれかの対応を行ってください:
- ダウンストリームの
xcom_pullを変更し、最後のelementを取得するようにする。 - ステートメントを
;で区切った単一の文字列として渡し、split_statements=Trueを設定する。デフォルトのreturn_last=Trueのままであれば、operatorは最後のステートメントの行のみをプッシュするため、pluginと同じ挙動になります。
parameters を使って行のリストを挿入していた ClickHouseOperator は、
SQLInsertRowsOperator になります。
columns を渡してください。指定しない場合、operator は SQLAlchemy 経由でテーブルを参照します。
カラム型を保持する
with_column_types=True は (rows, [(name, type), ...]) を返していました。これはハンドラーを使って再現できます。clickhouse-connect のカーソルは、cursor.description で ClickHouse の型名を返します:
ステップ5: ClickHouseHook.execute から DbApiHook のメソッドへ
プラグインの hook は、clickhouse_driver.Client.execute をそのまま反映した execute という単一のメソッドのみを公開していました。プロバイダの hook は DbApiHook であるため、他のすべての SQL プロバイダと同じ標準メソッド、すなわち run、get_records、get_first、get_pandas_df、get_df、insert_rows、test_connection が利用できます。コンストラクタ引数の clickhouse_conn_id と database は変更ありません。
fetch_all_handler およびその他のハンドラーは airflow.providers.common.sql.hooks.handlers から import します。
プラグイン時代のコードで最も多く見られる hook の書き方は bulk insert です。これは単純な run 呼び出しには置き換えられません。DB-API のカーソルが行を SQL 文字列へ整形しようとしてしまうためです。代わりに ネイティブ insert を使用してください:
変更前:
bulk_insert_rows には column_names が必要です。batch_size は任意で、非常に大きな入力に対してメモリ使用量を抑えます。汎用の insert_rows(table, rows, target_fields=[...], executemany=True) も最終的にはネイティブ insert になりますが、それは executemany=True の場合のみで、デフォルトでは 1 行ごとに 1 回の HTTP リクエストを送信します。
DB-API が提供する範囲でカバーできない操作については、get_client() が Airflow connection の設定を反映した生の clickhouse-connect client を返します。これは、plugin が公開していた clickhouse-driver 固有の argument すべてに代わるものです:
ステップ 6: ClickHouseSensor から SqlSensor へ
置き換えのなかで、callable への入力が変わるのはこれだけです。
plugin は result set 全体を渡していたため、既存のセンサーのコードはそこに索引でアクセスしています。移行時には、この索引アクセスを削除してください。
変更前:
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": ...}は引き続き動作します。 プロバイダーのフックはschemaをdatabaseの別名として扱います。 ClickHouseDbApiHookはClickHouseHookになります。コンストラクター引数schemaは引き続き別名として 受け付けられます。databaseがClickHouseネイティブの表記であり、両方が指定された場合はこちらが優先されます。- 接続には、ステップ2でのポートおよびextrasの変更が引き続き必要です。ラッパーもネイティブ プロトコルを使用していたためです。
確認すべき動作の違い
コードのコンパイルが通った後でも、実行時の動作が異なる点がいくつかあります。 INSERT タスクの XCom 値。 プラグインはclickhouse-driver が返した値をそのまま push しており、パラメータ付きの VALUES 挿入では挿入された行数が返されていました。一方 provider は、行を返さないステートメントに対しては空の結果セットを push します。XCom から行数を読み取っている下流のタスクは、後続の SELECT count() を実行するなど、別の方法で行数を取得する必要があります。
セッション設定と SET ステートメント。 どちらのパッケージも複数ステートメントのリストを単一の接続上で実行し、clickhouse-connect はデフォルトでクライアントごとにセッションを作成するため、リストの先頭付近にある SET ステートメントは後続のステートメントにも引き続き適用されます。それでも session_settings の使用を推奨します。明示的でテンプレート化も可能であり、サーバーや中間の proxy がセッションを保持するかどうかに関わらず同じように動作します。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 の依存関係として引き続きインストールされます。
型マッピング。 どちらのドライバーもネイティブな Python の型を返しますが、これらは別々のコードベースです。タイムゾーン付きの DateTime64、Decimal、UUID、Nullable カラム、ネストした Array や Map の値について、厳密な型に依存しているタスクを確認してください。特に結果が XCom に push されて下流で利用される場合は注意が必要です。
system.query_log でのクエリの識別。 クエリは HTTP インターフェイス経由で届くようになったため、interface = 1 ではなく interface = 2 として記録されます。また system.query_log の http_user_agent カラムには、Airflow と provider のバージョン、および client_name extra が設定されていればその値が含まれます。ネイティブプロトコルや clickhouse-driver のクライアント名でフィルタしていた監視は更新が必要です。行を返さない SELECT では 2 つ目のエントリが生成されます。これは DB-API カーソルがカラムのメタデータを取得するために SELECT * FROM (...) LIMIT 0 を実行するためです。
HTTP 経由のタイムアウト。 send_receive_timeout は HTTP の読み取りタイムアウトになり、worker と ClickHouse の間にある proxy やロードバランサーはリクエストに独自のアイドルタイムアウトを適用します。ネイティブプロトコル上で数分間実行されていたステートメントでは、これらの上限を引き上げる必要があるかもしれません。
接続の扱い。 hook は run または get_client の呼び出しごとに clickhouse-connect クライアントを作成します。これは、プラグインが execute ごとに新しいネイティブ接続を開いていたのと同じ考え方です。クライアントはプロセス全体で HTTP 接続プールを共有するため、get_client() から得たクライアントで close() を呼ぶことは必須ではなく、あくまで良い作法にすぎません。プールはタスクのプロセスが終了した時点で解放されます。クライアントはコンテキストマネージャーなので、with hook.get_client() as client: と書くのが最もすっきりします。
チェックリスト
- Airflow が 2.11 以降であること。
- HTTP ポートに worker から到達できること。HTTP エンドポイント用の TLS 証明書が有効であること。
- すべての ClickHouse connection について、type が
clickhouse、ポートが8123または8443、extras がステップ 2 に従って変換済みで、airflow connections testによる検証が済んでいること。 - ステップ 3 に従って import を置き換え、
common.sqlの wrapper からClickHouseプレフィックスを削除済みであること。 - operator と sensor の
clickhouse_conn_idをconn_idにリネームし、plugin のデフォルトに依存していたすべての task にconn_idを設定済みであること。 settings=をhook_params={"session_settings": ...}に移動済みであること。hook.execute("INSERT ... VALUES", rows)をbulk_insert_rowsに、ClickHouseOperator(parameters=rows)をSQLInsertRowsOperatorに置き換え済みであること。- sensor の callable を
result[0][0]ではなく cell の値そのものを参照するよう調整済みであること。 with_column_types、external_tables、columnar、query_id、types_checkの使用箇所をハンドラーまたはget_client()を使って書き換え済みであること。- 複数 statement の task および INSERT task が返す XCom を利用する後続のコンシューマーを見直し済みであること。
clickhouse_driverの exception をキャッチしているコードを更新済みであること。airflow-clickhouse-pluginとclickhouse-driverをアンインストール済みであること。