SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
4.1 KB · 115 lines typescript
Raw Blame History
1import { createClient, type ClickHouseClient } from "@clickhouse/client";23/**4 * ClickHouse: append-only analytics store — crawl history, extracted observations, page changes,5 * connector metrics and daily time series. Postgres remains the transactional source of truth.6 */7let _ch: ClickHouseClient | null = null;8export function getClickHouse(): ClickHouseClient {9  if (!_ch) {10    _ch = createClient({11      url: process.env.CLICKHOUSE_URL ?? "http://127.0.0.1:8123",12      database: process.env.CLICKHOUSE_DB ?? "dci",13      username: process.env.CLICKHOUSE_USER ?? "default",14      password: process.env.CLICKHOUSE_PASSWORD ?? "",15      request_timeout: 30_000,16      clickhouse_settings: { async_insert: 1, wait_for_async_insert: 0 },17    });18  }19  return _ch;20}2122export const CLICKHOUSE_DDL: string[] = [23  `CREATE TABLE IF NOT EXISTS crawl_log (24    ts DateTime64(3) DEFAULT now64(3),25    connector_id LowCardinality(String),26    source_id LowCardinality(String),27    document_id String,28    url String,29    fetch_level UInt8,30    fetcher LowCardinality(String),31    status_code UInt16,32    duration_ms UInt32,33    bytes UInt32,34    changed UInt8,35    not_modified UInt8,36    error LowCardinality(String),37    credits Float32,38    run_id String39  ) ENGINE = MergeTree ORDER BY (connector_id, ts) TTL toDateTime(ts) + INTERVAL 400 DAY`,40  `CREATE TABLE IF NOT EXISTS observations (41    ts DateTime64(3) DEFAULT now64(3),42    connector_id LowCardinality(String),43    source_id LowCardinality(String),44    document_id String,45    entity_type LowCardinality(String),46    entity_key String,47    entity_id String,48    field LowCardinality(String),49    value String,50    value_num Nullable(Float64),51    confidence LowCardinality(String),52    is_estimate UInt8,53    method LowCardinality(String),54    extractor_version LowCardinality(String),55    run_id String56  ) ENGINE = MergeTree ORDER BY (entity_type, entity_key, field, ts)`,57  `CREATE TABLE IF NOT EXISTS page_changes (58    ts DateTime64(3) DEFAULT now64(3),59    connector_id LowCardinality(String),60    document_id String,61    url String,62    old_hash String,63    new_hash String,64    added_lines UInt32,65    removed_lines UInt32,66    ratio Float32,67    significance UInt8,68    fields Array(String),69    event_types Array(String)70  ) ENGINE = MergeTree ORDER BY (connector_id, ts)`,71  `CREATE TABLE IF NOT EXISTS metrics (72    ts DateTime DEFAULT now(),73    metric LowCardinality(String),74    dim LowCardinality(String),75    value Float6476  ) ENGINE = MergeTree ORDER BY (metric, dim, ts) TTL ts + INTERVAL 400 DAY`,77  `CREATE TABLE IF NOT EXISTS entity_daily (78    day Date,79    metric LowCardinality(String),80    dim LowCardinality(String),81    value Float6482  ) ENGINE = ReplacingMergeTree ORDER BY (metric, dim, day)`,83  `CREATE TABLE IF NOT EXISTS api_requests (84    ts DateTime DEFAULT now(),85    route LowCardinality(String),86    status UInt16,87    duration_ms UInt32,88    cached UInt889  ) ENGINE = MergeTree ORDER BY (route, ts) TTL ts + INTERVAL 90 DAY`,90];9192export async function ensureClickHouse(): Promise<void> {93  const db = process.env.CLICKHOUSE_DB ?? "dci";94  const admin = createClient({ url: process.env.CLICKHOUSE_URL ?? "http://127.0.0.1:8123", username: process.env.CLICKHOUSE_USER ?? "default", password: process.env.CLICKHOUSE_PASSWORD ?? "" });95  await admin.command({ query: `CREATE DATABASE IF NOT EXISTS ${db}` });96  await admin.close();97  const ch = getClickHouse();98  for (const ddl of CLICKHOUSE_DDL) await ch.command({ query: ddl });99}100101export async function chInsert(table: string, rows: Record<string, unknown>[]): Promise<void> {102  if (!rows.length) return;103  try {104    await getClickHouse().insert({ table, values: rows, format: "JSONEachRow" });105  } catch (e) {106    if (process.env.DCI_CLICKHOUSE_OPTIONAL !== "0") console.warn(`[clickhouse] insert ${table} failed: ${(e as Error).message}`);107    else throw e;108  }109}110111export async function chQuery<T = Record<string, unknown>>(query: string, params: Record<string, unknown> = {}): Promise<T[]> {112  const rs = await getClickHouse().query({ query, query_params: params, format: "JSONEachRow" });113  return (await rs.json()) as T[];114}115