Skip to content

ClickHouse Cloud で実行可能 UDF が一般提供 (GA) となりました

TL;DR 実行可能 UDF が、AWS、GCP、Azure 上の ClickHouse Cloud で一般提供 (GA) となりました。コンパイル済みの Rust/Go/C++/JavaScript 向けの Native ランタイム、ネットワークアクセス、メモリ制限、deterministic フラグ、Cloud API および Terraform のサポート、26.6 以降でのクエリ単位の UDF メトリクスに対応しています。以下では、これらを使用して ClickHouse 内で LLM トークンをカウントして価格を算出し、回答の評価を行います。

注: 付属のコードはこちらにあります - github.com/ClickHouse/llm-token-udf。

4か月前、ClickHouse Cloud で実行可能 UDF のパブリックベータを開始しました。Python で関数を記述し、zip としてアップロードすれば、組み込み関数と同じように SQL から呼び出せます。ベータ期間中に寄せられた要望は、非常に一貫していました。Python ではなくコンパイル済みコードをアップロードしたい、Azure で UDF を使いたい、サポートチケットを開かずにネットワークアクセスを利用したい、そして UDF がクラスターに与えている影響を把握したい、というものです。

本日、実行可能 UDF が ClickHouse Cloud 上で一般提供 (GA) となったことをお知らせします(AWS、GCP、Azure で利用可能)。ベータ発表以降、以下の機能を追加しました。

  • Native ランタイム - Python スクリプトの代わりに、事前にコンパイルした静的リンクのバイナリをアップロードできます。現在は Rust、Go、C++、JavaScript(Bun でコンパイル )をサポートしています。
  • ネットワークアクセス - ベータ記事で紹介したプライベートプレビュー機能が、すべての組織で利用可能になりました。UDF からパブリックエンドポイントへのアウトバウンド呼び出しを行えます。
  • Azure - 3大クラウドすべてで UDF を実行できるようになりました。
  • プロセスごとのメモリ制限、および UDF を呼び出すクエリをクエリキャッシュから応答可能にする deterministic フラグ。
  • Cloud API における 12 個の UDF エンドポイント、および Terraform プロバイダー(3.24.0 以降)でのサービス単位のバージョン固定に対応した clickhouse_udf / clickhouse_udf_attachment リソース。
  • ClickHouse 26.6 以降における 8 個の ProfileEvents カウンターと 2 個の非同期メトリクス。これにより、クエリの UDF コスト(実経過時間、プール待機時間、CPU、メモリ、パイプ経由のバイト数)が他のすべての項目と並んで system.query_log に表示されます。

UDF は個別に請求されません。サービスの Pod 内で実行され、クエリと同じ CPU とメモリを消費します。

ベータの記事では、PyTorch のオートエンコーダーを使用して約 60 億件の株式取引をスコアリングしました。今回は、多くのユーザーから UDF の用途として伺っていた実態に近いもの、すなわち会計処理を取り上げます。具体的には、ClickHouse にすでにログとして記録しているプロンプトを使用して、LLM の請求額の内訳を把握します。UDF、SQL、Terraform の完全なソースコードは github.com/ClickHouse/llm-token-udf にあります。

なぜトークンなのか?

LLM 上で何らかのシステムを運用しているなら、プロンプト、完了テキスト、モデル、レイテンシ、そしておそらくトレース ID など、呼び出し内容をどこかに記録しているはずです。そのログの保存先として、ClickStack、Langfuse、または GenAI セマンティック規約を使用する OpenTelemetry コレクターを経由した ClickHouse の採用が増加しています。

しかし、確実に手元にないのがトークン数です。プロバイダーはほとんどのレスポンスで usage を返しますが、ストリーミングレスポンスでは返さない場合があり(返すもの、返さないもの、フラグ付きでのみ返すものがあります)、あらゆるプロキシを経由するわけでもなく、セルフホストされたモデルからも返されません。後述の合成データセットでは、スパンの約 30% が usage なしで到着します。これは、ユーザーにレスポンスをストリーミング配信する場合に概ね見られる傾向です。また、usage が存在する場合でも、それは合計値に過ぎません。サポートボットの入力コストの 41% が 3 月以降変更されていないシステムプロンプトによるものであることや、先週の火曜日に検索レイヤーが 2 倍のチャンクを返し始めたことまでは把握できません。

トークンのカウント自体はライブラリ呼び出し(Python では tiktoken、Rust では tiktoken-rs)なので、カウントすること自体が課題なのではありません。課題は、ライブラリがコード内に存在し、プロンプトがテーブル内に存在する点にあります。プロンプトをエクスポートし、ノートブックでカウントして結果を再度ロードすることも可能ですが、これでは処理が遅く、データが陳腐化し、管理すべきパイプラインがもう 1 つ増えることになります。arrayFold を使って 20 万語彙に対するバイトペアエンコーディングを実装することもできますが(推奨しません)、あるいはライブラリをデータのすぐそばに配置することもできます。

構築したもの

2 つの UDF といくつかの SQL を作成しました。

  • count_tokens(model, text) -> UInt32 - Native ランタイム上で動作する Rust バイナリ。決定論的(Deterministic)かつ CPU バウンドで、マテリアライズドビューにより挿入時にスパンあたり 4 回呼び出されます。
  • judge_response(prompt, completion) -> Tuple(score, verdict, reason) - ネットワークアクセスを持つ Python UDF。Claude に完了テキストのサンプルの評価を依頼し、tool-use を介して出力を制約します。リフレッシャブルマテリアライズドビューによって 10 分ごとに呼び出されます。

その他はすべて通常の ClickHouse の機能です。モデル価格の辞書(JSON ファイルの取得は UDF の仕事ではないため、url() を使用して公開価格リストから取得)、2 つのマテリアライズドビュー、そしてクエリです。

┌───────────────────────────┐
│  llm_spans                │     ← prompts, completions, model, usage (often NULL)
└──────────────┬────────────┘
               │ INSERT
               ▼
┌───────────────────────────┐
│  llm_span_tokens_mv       │     ← fires on every INSERT
│  (calls count_tokens ×4)  │
└──────────────┬────────────┘
               │
               ▼
┌───────────────────────────┐     ┌──────────────────────────┐
│  llm_span_tokens          │ ⟵─┤  model_prices (dict)     │ ← url() + refreshable MV, daily
│  system / context / user /│     └──────────────────────────┘
│  completion token counts  │
└───────────────────────────┘
               ▲
               │ every 10 min, 2% sample
┌──────────────┴────────────┐
│  llm_evals_mv             │     ← refreshable MV, APPEND
│  (calls judge_response)   │        over the network
└───────────────────────────┘

トークナイザーを Native UDF として実装する

Native ランタイム(どの UDF でも選択可能な I/O フォーマットの 1 つである ClickHouse の Native フォーマットと混同しないでください)は、amd64/ と arm64/ の 2 つのフォルダーを含む zip を受け取ります。各フォルダーには、main という名前の静的リンクされた Linux 実行ファイルと、実行時に併せて必要なデータファイルが含まれます。ClickHouse Cloud は両方のアーキテクチャで稼働しているため、両方が必須となります。バイナリは引数なしで実行され、自動的にインストールされるものは何もないため、必要なものはすべてコンパイルして組み込む必要があります。

count_tokens は、tiktoken-rs をラップした約 100 行の Rust コードです。大部分はワイヤプロトコルの処理であり、これは Python UDF がやり取りするものと同じですが、バイナリ形式で処理されます。

// Deploy as: executable_pool, runtime = Native, format = RowBinary,
// send_chunk_header = true, deterministic = true.
// Arguments: (model String, text String) -> UInt32
fn main() {
    let mut tok = Tokenizers::load();           // reads models.json from the binary's directory
    let mut stdin = BufReader::with_capacity(1 << 20, io::stdin().lock());
    let mut stdout = BufWriter::with_capacity(1 << 20, io::stdout().lock());
    let (mut header, mut model, mut text) = (String::new(), Vec::new(), Vec::new());

    loop {
        // 1. chunk header: row count as decimal text + '\n' (send_chunk_header)
        header.clear();
        match stdin.read_line(&mut header) {
            Ok(0) => return,                    // pipe closed; pool process exits cleanly
            Ok(_) => {}
            Err(e) => die(&format!("reading chunk header: {e}")),
        }
        let rows: usize = header.trim().parse().unwrap_or_else(|_| die("bad chunk header"));

        // 2. N RowBinary rows: each String is a LEB128 length + raw bytes
        for i in 0..rows {
            read_string(&mut stdin, &mut model)
                .and_then(|_| read_string(&mut stdin, &mut text))
                .unwrap_or_else(|e| die(&format!("row {i}: {e}")));
            let n = tok.count(&String::from_utf8_lossy(&model), &String::from_utf8_lossy(&text));
            stdout.write_all(&n.to_le_bytes()).unwrap_or_else(|_| process::exit(1));
        }
        // 3. N UInt32 answers, flushed once per chunk
        stdout.flush().unwrap_or_else(|_| process::exit(1));
    }
}

ベータ版のデモで使用した TabSeparated フォーマットではなく、RowBinary を採用しました。14 個の数値特徴量であれば TSV で十分でしたが、プロンプトにはパイプの両側でエスケープが必要となるタブ、改行、バックスラッシュが多く含まれます。これに対して RowBinary はバイナリセーフであり、ClickHouse が送信時にテキストをフォーマットしたり受信時にパースしたりする必要がありません。検証では、UDF を TabSeparatedRaw から RowBinary に移行するだけで約 20% の性能向上が得られました。

チャンクヘッダーはプロトコルのもう 1 つの重要な要素です。send_chunk_header を有効にすると、ClickHouse は各チャンクの前に処理行数を書き込むため、プロセスは正確にその行数だけを読み取り、処理して、1 回だけフラッシュします。これがないと、出力をバッファリングする UDF はブロックがいつ終了するかを知る方法がありません。ベータ期間中、128 行ごとにフラッシュする UDF をデバッグしましたが、その UDF は 128 の倍数ではないサイズのすべてのブロック(つまり、ほぼすべてのクエリの最後のブロック)で、読み取りタイムアウトが発生するまでハングアップしていました。UDF が行ごとにフラッシュしない場合は、この設定を有効にしてください。

モデルとエンコーディングのマッピング(gpt-4o → o200k_base、gpt-4 → cl100k_base など)は、バイナリと同じ階層にある models.json で管理されています。そのため、新しいモデルが登場した際も再コンパイルは不要で、テキストファイルを編集して新しいバージョンをアップロードするだけで対応できます。レビューの過程で判明した注意点として、データファイルはバイナリと同じ階層に配置されますが、プロセスはそのディレクトリ内で起動されるわけではありません(バンドルは /scripts 配下に配置され、作業ディレクトリは / となります)。そのため、バイナリは作業ディレクトリではなく自身のパスを基準に models.json を解決し、ファイルが存在しない場合は警告なしに処理を継続するのではなく、エラーを出して終了するように設計されています。認識できないモデルは o200k_base にフォールバックします。これは OpenAI 以外のトークナイザーに対しては近似値となりますが、コスト配分を把握する目的には十分であり、プロバイダー独自の usage が存在する場合はそちらが常に優先されます(詳細は後述)。

ラップトップから両方のアーキテクチャ向けにビルドするには、musl ターゲットを指定した 2 つの cargo コマンドを実行します。

cargo build --release --target x86_64-unknown-linux-musl
cargo build --release --target aarch64-unknown-linux-musl
# -> count_tokens.zip: amd64/main, amd64/models.json, arm64/main, arm64/models.json

バイナリはそれぞれ約 7MB、zip ファイルは 6.6MB です。

デプロイ

デプロイ画面はベータ版と同じアップロード画面ですが、いくつか新しいフィールドが追加されています。

フィールド値
Namecount_tokens
Typeexecutable_pool
Runtime typeNative
FormatRowBinary
Send chunk headertrue
Deterministictrue
Memory limit512 MiB
Pool size4
Argumentsmodel String, text String
Return typeUInt32

これらのうち 2 つがベータ以降に追加された新しい項目です。「Deterministic」は、関数が同じ入力に対して同じ出力を返すことを ClickHouse に伝えます。これは、クエリキャッシュが結果を保存する前に確認する必要がある情報です。ベータ期間中はすべての Cloud UDF が非決定論的(non-deterministic)として扱われていたため、UDF を呼び出すクエリで use_query_cache = 1 を指定すると QUERY_CACHE_USED_WITH_NONDETERMINISTIC_FUNCTIONS で失敗していました。このフラグを設定すると、ダッシュボードクエリの 2 回目の実行は、UDF を一切呼び出さずに 0 ms でキャッシュから結果を返します。

SELECT feature, sum(count_tokens(model, system_prompt)) AS system_tokens
FROM llm_spans
WHERE toDate(ts) = '2026-10-01'    -- a literal on purpose: today() is itself non-deterministic
GROUP BY feature
SETTINGS use_query_cache = 1;
┌─query_duration_ms─┬─QueryCacheHits─┬─QueryCacheMisses─┬─udf_invocations─┐
│               209 │              0 │                1 │              14 │   ← first run
│                 0 │              1 │                0 │               0 │   ← second run
└───────────────────┴────────────────┴──────────────────┴─────────────────┘

「Memory limit」は各サンドボックスプロセスが使用できるメモリの上限を設定します。デフォルトは 4 GiB です。この制限は実メモリ(resident memory)ではなく仮想アドレス空間に適用されますが、一部のランタイムは実際に使用するよりもはるかに多くのアドレス空間を確保するため、この点は重要です。作成した Rust バイナリはプロセスあたり仮想 43 MiB / 実メモリ 39 MiB であり、Python 版は 103 / 93 MiB、Go 版は 29 MiB を使用しながら約 1.2 GiB のアドレス空間を予約します。つまり、ワークロードではなくランタイムに合わせて制限サイズを設定し、Go の場合は十分な余裕を持たせる必要があります。

取り込みパイプラインへの組み込み

マテリアライズドビューから count_tokens を呼び出すことで、トークナイザーが挿入時にスパンあたり正確に 1 回だけ実行されるようにします。

CREATE MATERIALIZED VIEW llm_span_tokens_mv TO llm_span_tokens AS
SELECT
    ts, trace_id, span_id, service, feature, customer_id, model,
    count_tokens(model, system_prompt)     AS system_tokens,
    count_tokens(model, retrieved_context) AS context_tokens,
    count_tokens(model, user_prompt)       AS user_tokens,
    count_tokens(model, completion)        AS completion_tokens,
    provider_input_tokens,
    provider_output_tokens,
    sipHash64(system_prompt)               AS system_prompt_hash
FROM llm_spans;

INSERT INTO llm_spans が実行されるたびにこのビューがトリガーされ、スパンごとに 4 つの要素をカウントして、その結果を llm_span_tokens に格納します。それ以降のすべてのクエリは、整数値に対する集計処理となります。

スループットの目安として、4 vCPU、プールサイズ 4 の環境において、10 万件の合成スパン(351 MiB のプロンプトテキスト、6,900 万トークン)がこのビューを 5.2 秒で通過しました。これは毎秒約 1,300 万トークン、または毎秒約 1 万 9,000 スパンに相当します。

価格の算出

トークン数にトークン単価を掛けることで、金額を算出できます。LiteLLM は公開価格リストを JSON ファイルとして管理しています。ファイルの取得は UDF ではなく url() の役割であるため、その取得処理をリフレッシャブルマテリアライズドビューでラップし、毎日再取得して辞書に投入します。

CREATE MATERIALIZED VIEW model_prices_mv
REFRESH EVERY 1 DAY
TO model_prices
AS
SELECT
    kv.1                                                   AS model,
    JSONExtractString(kv.2, 'litellm_provider')            AS provider,
    JSONExtractFloat(kv.2, 'input_cost_per_token')         AS input_cost_per_token,
    JSONExtractFloat(kv.2, 'output_cost_per_token')        AS output_cost_per_token,
    JSONExtractFloat(kv.2, 'cache_read_input_token_cost')  AS cache_read_input_token_cost,
    now()                                                  AS updated_at
FROM
(
    SELECT arrayJoin(JSONExtractKeysAndValuesRaw(json)) AS kv
    FROM url('https://raw.githubusercontent.com/BerriAI/litellm/main/model_prices_and_context_window.json',
             JSONAsString, 'json String')
)
WHERE JSONHas(kv.2, 'input_cost_per_token');

これにより、1 日 1 回更新され、dictGet で参照可能な 3,683 個のモデルの価格情報が得られます。コスト計算クエリでは、プロバイダーの usage が存在する場合はそれを使用し、存在しない場合は独自にカウントした値を使用します。

SELECT
    toDate(ts) AS day,
    feature,
    round(sum(coalesce(provider_input_tokens, system_tokens + context_tokens + user_tokens)
            * dictGet('model_prices_dict', 'input_cost_per_token', model)
          + coalesce(provider_output_tokens, completion_tokens)
            * dictGet('model_prices_dict', 'output_cost_per_token', model)), 2)  AS est_cost_usd,
    countIf(provider_input_tokens IS NULL)                         AS spans_filled_by_udf,
    count()                                                        AS spans
FROM llm_span_tokens
GROUP BY day, feature
ORDER BY day DESC, est_cost_usd DESC;
┌────────day─┬─feature───────┬─est_cost_usd─┬─spans_filled_by_udf─┬─spans─┐
│ 2026-09-30 │ sql_assistant │         3.25 │                 558 │  1781 │
│ 2026-09-30 │ support_bot   │         3.24 │                 539 │  1759 │
│ 2026-09-30 │ docs_search   │         3.06 │                 535 │  1796 │
│ 2026-09-30 │ summarizer    │         2.93 │                 527 │  1818 │
└────────────┴───────────────┴──────────────┴─────────────────────┴───────┘

プロバイダーが使用量を報告しているケースでは、独自に算出した結果と照合できます。provider_input_tokens と算出された 3 つの要素のカウント値の合計との差分は、チャットテンプレートのオーバーヘッド(ロールマーカーやメッセージの枠組み)であり、特定のモデルであれば安定して十数トークン程度になるはずです。これに乖離が生じた場合、プロバイダーがテンプレートを変更したか、そのモデルに対する models.json の設定が誤っているかのいずれかであり、どちらの場合も 1 行の quantile クエリですぐに調査できます。

各要素からわかること

システムプロンプトは呼び出しごとに送信される静的なテキストであり、現在プロバイダーはキャッシュされたプレフィックストークンを定価の 50〜90% 引きで請求します(前述の cache_read_input_token_cost カラムがこれに該当します)。したがって、「入力コストのうちシステムプロンプトが占める割合はどれくらいか?」という問いは、「プロンプトキャッシングによってどれだけのコストを削減できるか?」という問いと同じ意味を持ちます。

SELECT
    feature,
    sum(system_tokens)                                              AS sys_tokens,
    sum(system_tokens + context_tokens + user_tokens)               AS input_tokens,
    round(sys_tokens / input_tokens * 100, 1)                       AS system_pct,
    round(sum(system_tokens
              * (dictGet('model_prices_dict', 'input_cost_per_token', model)
                 - dictGet('model_prices_dict', 'cache_read_input_token_cost', model))), 2)
                                                                    AS usd_saved_if_cached
FROM llm_span_tokens
WHERE ts >= now() - INTERVAL 7 DAY
GROUP BY feature
ORDER BY usd_saved_if_cached DESC;
┌─feature───────┬─sys_tokens─┬─input_tokens─┬─system_pct─┬─usd_saved_if_cached─┐
│ support_bot   │    3483276 │      8461730 │       41.2 │                4.04 │
│ sql_assistant │    3237735 │      8382894 │       38.6 │                3.78 │
│ docs_search   │    2213208 │      7287850 │       30.4 │                2.61 │
│ summarizer    │    1552014 │      6999309 │       22.2 │                1.81 │
└───────────────┴────────────┴──────────────┴────────────┴─────────────────────┘

データセットが小規模(2 週間で 10 万スパン)であるため金額は小さくなっていますが、実際の運用規模に応じてスケールさせて考えてみてください。重要なのは、このクエリが usage の合計値だけでは決して作成できないという点です。内訳となる要素が必要であり、その要素はトークナイザーから取得されます。

ネットワーク経由でのサンプルの評価

コストの把握は会計処理の半分に過ぎず、生成された回答の品質を検証することがもう半分となります。一般的な手法は、別のモデルにサンプルの評価を依頼することです。これはデータベース内部から LLM API へのネットワーク呼び出しを行う処理であり、ネットワーク対応 UDF がまさにこの用途に適しています。

judge_response は、プロンプトと完了テキストを Claude に送信する約 120 行の Python UDF です。tool-use スキーマを使用しており、その verdict フィールドは pass、partial、fail の enum、score は 1 から 5 までの整数値となっています。モデルにはこのツールの呼び出しが義務付けられているため、レスポンスは常にパース可能であり、verdict は必ずこれら 3 つの値のいずれかになります。

# abridged - full file in the repo
JUDGE_TOOL = {
    "name": "report_verdict",
    "description": "Grade how well the completion answers the prompt.",
    "input_schema": {
        "type": "object",
        "properties": {
            "score":   {"type": "integer", "minimum": 1, "maximum": 5},
            "verdict": {"type": "string", "enum": ["pass", "partial", "fail"]},
            "reason":  {"type": "string", "maxLength": 200},
        },
        "required": ["score", "verdict", "reason"],
    },
}

def main() -> None:
    for line in sys.stdin:                       # JSONEachRow in ...
        if not line.strip():
            continue
        row = json.loads(line)
        score, verdict, reason = judge(row["prompt"], row["completion"])
        sys.stdout.write(json.dumps({"result": [score, verdict, reason]}) + "\n")   # ... JSONEachRow out
        sys.stdout.flush()

各行が HTTPS のラウンドトリップを伴う場合はスループットを重視する必要がなく、また JSON を使用するとタプルの戻り値を簡単に出力できるため、ここでは RowBinary ではなく JSONEachRow を採用しています。各プールプロセスは、その生存期間中に requests.Session を維持し、(prompt, completion) 単位で結果をキャッシュして再評価のコストをゼロにし、429 や 529 のエラーが発生した場合はバックオフを伴うリトライを実行します。また、無人で実行されるため、ソフトフェイルするように設計されており、不正な API 呼び出しが発生してもクエリを中断せず、(0, 'error', reason) を返します。

API キーについて: 現時点では UDF 向けのシークレットマネージャーが存在しないため(後述)、キーは SQL 上に露出させるのではなく、zip 内の config.json に含めて配置します。サンドボックスには独自のクラウド ID(インスタンスロール、メタデータエンドポイント、継承された環境変数など)が存在しないため、UDF が保持する認証情報はユーザーが付与したもののみとなります。

フィールド値
Namejudge_response
Typeexecutable_pool
Runtime typepython3.11
FormatJSONEachRow
Network accessenabled
Deterministicfalse
Max command exec time (s)30
Pool size4
Argumentsprompt String, completion String
Return typeTuple(UInt8, String, String)

評価ループ自体はリフレッシャブルマテリアライズドビューです。10 分ごとに直近 10 分間のスパンから決定論的に 2% のサンプルを抽出し、評価を実行して結果を追加します。

CREATE MATERIALIZED VIEW llm_evals_mv
REFRESH EVERY 10 MINUTE
APPEND TO llm_evals
AS
WITH judge_response(
        concat(system_prompt, '\n\n', retrieved_context, '\n\nUser: ', user_prompt),
        completion
     ) AS j
SELECT
    ts, trace_id, span_id, feature, model,
    j.1   AS score,
    j.2   AS verdict,
    j.3   AS reason,
    now() AS judged_at
FROM llm_spans
WHERE ts >= now() - INTERVAL 10 MINUTE
  AND cityHash64(trace_id) % 100 < 2;

ここにはスケジューラーや個別の評価サービスは存在せず、更新間隔を設定したビューがあるだけです。コスト集計で作成したような GROUP BY feature, model を実行するだけで、機能やモデルごとの失敗率を取得できます。さらに、llm_evals と llm_span_tokens は span_id を共有しているため、「どの高コストなプロンプトが同時に失敗しているか」を JOIN 処理で容易に特定できます。

実行コスト

ベータ版における最大の課題はオブザーバビリティでした。UDF は独立したプロセスで実行されるため、system.query_log にはクエリ自体の CPU やメモリのみが表示され、子プロセスに関する情報は一切出力されませんでした。そのため、UDF とクエリのどちらがボトルネックになっているのかをユーザーから問い合わせられた際、システムテーブルから回答することができませんでした。

ClickHouse 26.6 では、実行可能 UDF 向けに 8 つの ProfileEvents が追加され、すべての Cloud サービスの system.query_log、system.events、system.metric_log に表示されるようになりました。

ProfileEvent測定対象
ExecutableUserDefinedFunctionInvocationsUDF の呼び出し回数(チャンクごとに 1 回)
ExecutableUserDefinedFunctionElapsedMicrosecondsUDF 呼び出しに費やされた実経過時間
ExecutableUserDefinedFunctionPoolWaitMicroseconds空きプールプロセスを待機した時間
ExecutableUserDefinedFunctionUserTimeMicroseconds子プロセスによって消費されたユーザーモード CPU
ExecutableUserDefinedFunctionSystemTimeMicroseconds子プロセスによって消費されたカーネルモード CPU
ExecutableUserDefinedFunctionPeakMemoryByteSecondsプロセスごとのピークメモリの実経過時間に対する積分値
ExecutableUserDefinedFunctionInputBytes子プロセスの stdin に書き込まれたバイト数
ExecutableUserDefinedFunctionOutputBytes子プロセスの stdout から読み取られたバイト数

また、2 つの非同期メトリクスである ExecutableUserDefinedFunctionProcesses と ExecutableUserDefinedFunctionMemoryResidentBytes により、サーバー上で現在稼働している UDF プロセスの数と、それらが保持している実メモリ量(プロセスごとに合計されるため、上限値となります)がレポートされます。

これらが導入されたことで、UDF をどの言語で記述すべきかという疑問は、query_log をクエリするだけで解決できるようになりました。Native ランタイムの Rust、tiktoken を使用した Python、そして純粋な Go 実装のトークナイザーを使用した Go の 3 つの言語で count_tokens を作成し、マテリアライズドビューのワークロード(10 万スパン、それぞれ 4 回のカウント、351 MiB のテキスト)をそれぞれに流してテストしました。

SELECT
    extract(query, 'AS (\\w+)_tokens')                                            AS impl,
    query_duration_ms,
    ProfileEvents['ExecutableUserDefinedFunctionInvocations']                     AS invocations,
    round(ProfileEvents['ExecutableUserDefinedFunctionElapsedMicroseconds'] / 1e6, 2) AS udf_wall_s,
    round((ProfileEvents['ExecutableUserDefinedFunctionUserTimeMicroseconds']
         + ProfileEvents['ExecutableUserDefinedFunctionSystemTimeMicroseconds']) / 1e6, 2) AS udf_cpu_s,
    formatReadableSize(ProfileEvents['ExecutableUserDefinedFunctionInputBytes'])  AS udf_in
FROM system.query_log
WHERE type = 'QueryFinish' AND ProfileEvents['ExecutableUserDefinedFunctionInvocations'] > 0
ORDER BY event_time_microseconds;
┌─impl─┬─query_duration_ms─┬─invocations─┬─udf_wall_s─┬─udf_cpu_s─┬─udf_in─────┐
│ rs   │              5161 │         144 │       7.05 │      6.94 │ 355.09 MiB │
│ py   │              7207 │         144 │      10.05 │      9.94 │ 355.09 MiB │
│ go   │             15028 │         144 │      20.97 │     26.14 │ 355.09 MiB │
└──────┴───────────────────┴─────────────┴────────────┴───────────┴────────────┘

この検証では、Rust は Python より 1.4 倍高速で、Go より 2.9 倍高速でした。これは当初想定していた順序とは異なっていました。Python が健闘しているのは、tiktoken のコア部分が Rust で実装されているためです。1.4 倍の差はインタープリターのループとパイプ処理によるものです。Go が遅い要因は、純粋な Go 実装のトークナイザーライブラリ自体が低速であるためであり、Native ランタイムを使用しても低速なライブラリそのものが高速化されるわけではありません。簡単に言えば、ランタイムはインタープリターや依存関係のインストールを排除し、手元にあるライブラリをそのまま実行できるようにしますが、ライブラリ自体が高速である必要があります(低レベルのネイティブコアを持たない、純粋に CPU バウンドな Python コードの場合は結果が異なります。デザインパートナーのあるユーザーは、文字列マッチング UDF を Python から Go に移行することで約 25 倍の高速化を測定しました)。

参考までに、算出されたカウント値は 3 つの実装すべてで完全に一致し、スポットチェックした 1,200 件の値において tiktoken 自体の結果とも一致しました。

Terraform

ここまでの内容はすべてコンソール上で設定しました。本番環境への導入にあたってはコードで管理したいと考えるのが一般的であり、Terraform プロバイダー(3.24.0 以降)ではそのための 2 つのリソースが利用可能になっています。clickhouse_udf は zip のハッシュが変更されるたびに新しいバージョンを発行してビルド完了を待機し、clickhouse_udf_attachment は特定のバージョンを特定のサービスに関連付けます。

resource "clickhouse_udf" "count_tokens" {
  function_name = "count_tokens"
  runtime       = "native"
  type          = "executable_pool"
  format        = "RowBinary"
  return_type   = "UInt32"
  arguments     = [{ name = "model", type = "String" }, { name = "text", type = "String" }]

  pool_size                  = 4
  send_chunk_header          = true
  max_command_execution_time = 10

  source_archive_path = "${path.module}/../udf/count_tokens_rs/count_tokens.zip"
  source_archive_hash = filebase64sha256("${path.module}/../udf/count_tokens_rs/count_tokens.zip")
}

# Dev follows every successful build.
resource "clickhouse_udf_attachment" "dev" {
  function_name = clickhouse_udf.count_tokens.function_name
  service_id    = var.dev_service_id
  version       = clickhouse_udf.count_tokens.version
}

# Prod stays where you pinned it.
resource "clickhouse_udf_attachment" "prod" {
  function_name = clickhouse_udf.count_tokens.function_name
  service_id    = var.prod_service_id
  version       = var.count_tokens_prod_version
}

これはコンソールでのバージョン管理の仕組みと同じです。バージョンはイミュータブル(不変)であり、新しい zip は新しいバージョンとなります。そしてバージョンはサービスごとにアタッチされるため、本番環境を v6 に維持したまま開発環境で v7 を実行したり、他のサービスに影響を与えることなく単一のサービスのみをロールバックしたりできます。各サービスが一度に保持できる関数のバージョンは 1 つだけです。executable_pool UDF の場合、長時間稼働しているプールプロセスはプールがリフレッシュされるまで古いバージョンを処理し続けます。そのため、コンソールのアタッチされたサービスの横には、内部で SYSTEM RELOAD FUNCTION を実行する「Reload UDF」アクションが用意されています。新しい設定項目のうち deterministic とメモリ制限の 2 つはプロバイダースキーマにまだ含まれていないため、現時点ではコンソールまたは API 経由で設定します。

同様のライフサイクル管理は Cloud API でも利用できます。アップロード URL を作成し、zip をプッシュして、関数または新しいバージョンを作成し、それをサービスにアタッチします。

ベータで得られた知見

4か月にわたるベータ期間中のサポートスレッドで得られた知見は、簡潔なリストに集約されます。

executable_pool を使用すること。通常の executable では、データのブロックごとに ClickHouse が新しいサンドボックスプロセスを起動します。これに対し、executable_pool は長時間稼働するプロセスのプールをブロックやクエリをまたいで再利用します。これにより処理が高速化され、モデルや語彙がメモリ上にウォーム状態で保持されるため、ブロックごとにプロセスを生成する方式では耐えられないような負荷にも耐えることができます。Cloud のワークロードにおいて、executable が適切な選択肢となったケースはまだ確認されていません。

チャンクの数を減らし、サイズを大きくすること。呼び出しには、パイプの両側で一定のセットアップコストが発生します。同じ 5,000 行のクエリでも、1 つのチャンクとして処理した場合は 165 ms でしたが、max_block_size = 1(つまり 5,000 回の呼び出し)とした場合は 1,544 ms かかりました。単一のクエリで ExecutableUserDefinedFunctionInvocations が数十万回に達している場合は、preferred_block_size_bytes の値を引き上げ(検証では 100000000 を使用)、再度カウンターを確認してください。

プールサイズはレプリカ単位であり、上限は max_threads となること。クエリが同時に使用するプールプロセスの数は最大でも max_threads までです。したがって、16 スレッドのクエリに対してプールサイズを 64 に設定しても、実質的には 16 のプールとして動作します。まずは 4 から開始し、ExecutableUserDefinedFunctionPoolWaitMicroseconds の値に応じて引き上げてください。その際、各プロセスが起動時に読み込んだデータのコピーをそれぞれ保持する点に留意してください。

N 行ごとではなく、チャンクごとにフラッシュすること。send_chunk_header を有効にし、指定された正確な行数を読み取ってください。前述の通りですが、これはベータ期間中のサポートスレッドで最も多く見られた問題でした。

エラー時は明示的に失敗させること。同梱されたファイルが見つからない、設定をパースできない、理解できない行を受信したなどの場合、UDF は stderr に 1 行出力してゼロ以外のステータスで終了すべきです。ClickHouse は stderr の内容をクエリのエラー情報として表示するため、3 つ後のダッシュボードで誤った数値を見て気づくのではなく、最初のクエリの時点で問題を把握できます。最初に作成した count_tokens では、models.json が見つからない場合に「ルールなし」として扱い、すべてのモデルを o200k_base としてトークナイズしていました。カウント値自体はもっともらしく見えていたため、gpt-4 の請求額の突合が合わなくなるまで誰も気づかなかったはずです。

UDF が不適切なツールとなる場面を見極めること。すべての行はパイプを介してプロセス境界を越え、その過程でシリアライズされるため、このコストをなくすことはできません。単純な行単位の処理ではこのオーバーヘッドが支配的になります。何もしない UDF を数十億行に対して実行した場合、どの言語であっても同等の組み込み関数より数倍遅くなります。UDF が真価を発揮するのは、SQL で表現できないロジック(トークナイザー、モデル、パーサー、API 呼び出しなど)を扱う場合であり、SQL で表現できるロジックを置き換えるためのものではありません。

サービスへの最初の UDF アタッチ時にはサービスの再起動が発生すること。最初の UDF をアタッチすると、サービスの Pod にヘルパーコンテナが追加されるため、ローリング再起動が発生します。2 つ目以降の UDF のアタッチでは再起動は発生しませんが、最後の UDF を削除する際には再度再起動が発生します。初回の UDF の関連付けは、メンテナンスウィンドウとして扱ってください。

今後のロードマップ

  • シークレット管理 - 最も多くの要望をいただいている機能です。API キーを zip 内に含めずに済むよう、UDF 向けの環境変数またはシークレット参照をサポートします。適切なアクセス分離を備えた Named Collection に UDF の入力を紐付ける設計を進めており、次の優先項目となっています。
  • 再ビルドを伴わない実行時設定の変更 - 現在、プールサイズ、タイムアウト、メモリ制限はバージョンに紐付けられています。新しいバージョンを発行することなく、clickhouse_udf_attachment を介してサービス単位でこれらをオーバーライドできるように分離作業を進めています。
  • Native UDF 向けのネットワークアクセス - 現在、Native ランタイムはコンピュート処理のみに対応しており、外部ネットワークへのアクセスは Python のみで利用可能です。
  • さらなるランタイムの追加 - WebAssembly UDF はオープンソースの ClickHouse に実験的機能として存在していますが、Cloud にはまだ導入されていません。より長期的なロードマップとして、WebAssembly や、各言語向けに最適化されたランタイムの提供を検討しています。

ぜひお試しください

実行可能 UDF は、AWS、GCP、Azure 上のすべての ClickHouse Cloud 組織で本日から一般提供 (GA) されており、追加の設定なしですぐに利用できます。コンソールの組織メニューから「User-defined functions」を開くか、ドキュメント、Cloud API リファレンス、または Terraform リソース を参照してください。

プロジェクト全体のコードは github.com/ClickHouse/llm-token-udf に公開されています。

llm-token-udf/
├── udf/
│   ├── count_tokens_rs/   # the Native UDF (Rust): source, build script, zip layout
│   ├── count_tokens/      # the same UDF in Go
│   ├── count_tokens_py/   # the same UDF in Python
│   └── judge_response/    # the network UDF (Python, Anthropic tool-use)
├── sql/                   # schema, MVs, prices, evals, and every query in this post
├── terraform/             # clickhouse_udf + dev/prod attachments
└── local/                 # XML config for running the same UDFs against open-source ClickHouse

Python UDF を Native ランタイムに移植して性能差を測定された方や、私たちが思いもよらなかったユースケースを構築された方は、ぜひお知らせください。

今すぐ始める

今すぐ ClickHouse Cloud を始めて、$300 のクレジットを受け取りましょう。30 日間の無料トライアル終了後は、従量課金プランに移行できます。ボリュームベースの割引について詳しくは お問い合わせ ください。詳細は 料金ページ をご覧ください。


この記事をシェア

  • Y Combinator icon
  • X icon
  • Bluesky icon
  • Facebook icon
  • LinkedIn icon

Subscribe to our newsletter

Stay informed on feature releases, product roadmap, support, and cloud offerings!

Follow us

XBlueskySlackGithubTelegramMeetupRSS