spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
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