Skip to main content
Ask your AI

ClickHouse

Server Source code Package Beta

Server-side event storage in ClickHouse. Raw walkerOS events are written in bulk as JSONEachRow inserts into a table you own, with a bounded retry made safe by an explicit insert_deduplication_token. The destination issues INSERT and nothing else: it never queries, never runs DDL, and never reads system.*, so the table's schema, sorting key, partitioning and retention stay entirely yours.

An insert into a table that does not exist fails with 60 UNKNOWN_TABLE, so read Create the table and Preconditions before sending the first event.

Where this fits

ClickHouse is a server destination in the walkerOS flow:

Receives events server-side from the collector, converts each one into a row of JSON-encoded payload columns, and inserts a flushed batch as a single JSONEachRow statement for analytics and reporting.

Installation​

npm install @walkeros/server-destination-clickhouse
import { startFlow } from '@walkeros/collector';
import { destinationClickHouse } from '@walkeros/server-destination-clickhouse';

await startFlow({
  destinations: {
    clickhouse: {
      code: destinationClickHouse,
      config: {
        settings: {
          url: 'https://clickhouse.example.com:8443',
          database: 'analytics',
          table: 'events',
        },
        credentials: {
          username: 'walkeros_writer',
          password: process.env.CLICKHOUSE_PASSWORD ?? '',
        },
        batch: { size: 5000, age: 5000, wait: 1000 },
      },
    },
  },
});

Configuration​

This destination uses the standard destination config wrapper (consent, data, env, id, ...). For the shared fields see destination configuration. Package-specific fields live under config.settings and are listed below.

Settings​

PropertyTypeDescriptionMore
url*stringClickHouse HTTP endpoint
databasestringDatabase holding the events table
tablestringTable every event is inserted into
maxRetriesintegerRetries added on top of the first attempt, so total attempts = 1 + maxRetries
clickhouseRecord<string, any>Raw @clickhouse/client options merged into the created client
clickhouseSettingsRecord<string, any>ClickHouse settings applied to every insert
* Required fields

Mapping​

This package does not define custom rule-level settings. For the standard rule fields (consent, condition, data, batch, name, policy) see mapping.

Examples

Initialization

Init creates a single ClickHouse client for the configured endpoint, database and credentials. The client is held on the resolved config and reused by every insert. Its request timeout is derived from the collector timeout so a retried insert cannot outlive the delivery it belongs to, and request compression is enabled.

Event
{
  "settings": {
    "url": "https://clickhouse.example.com:8443",
    "database": "analytics",
    "table": "events"
  },
  "credentials": {
    "username": "walkeros_writer",
    "password": "$secret.CLICKHOUSE_PASSWORD"
  }
}
Out
createClient({
  "url": "https://clickhouse.example.com:8443",
  "database": "analytics",
  "username": "walkeros_writer",
  "password": "$secret.CLICKHOUSE_PASSWORD",
  "request_timeout": 4812,
  "compression": {
    "request": true
  }
})

Page view

A page view is inserted as a single JSONEachRow row. The payload columns (data, context, globals, custom, user, nested, consent, source) are JSON-encoded strings and the timestamp is a UTC datetime literal, both produced by eventToRow.

Event
{
  "name": "page view",
  "data": {
    "domain": "www.example.com",
    "title": "walkerOS documentation",
    "referrer": "https://www.walkeros.io/",
    "search": "?foo=bar",
    "hash": "#hash",
    "id": "/docs/"
  },
  "context": {
    "dev": [
      "test",
      1
    ]
  },
  "globals": {
    "pagegroup": "docs"
  },
  "custom": {
    "completely": "random"
  },
  "user": {
    "id": "us3r",
    "device": "c00k13",
    "session": "s3ss10n"
  },
  "nested": [
    {
      "entity": "child",
      "data": {
        "is": "subordinated"
      }
    }
  ],
  "consent": {
    "functional": true
  },
  "id": "ca81a74de5cb7d41",
  "trigger": "load",
  "entity": "page",
  "action": "view",
  "timestamp": 1700000000123,
  "timing": 3.14,
  "source": {
    "count": 1,
    "trace": "0a1b2c3d4e5f60718293a4b5c6d7e8f9",
    "type": "collector",
    "schema": "4"
  }
}
Out
insert({
  "table": "events",
  "values": [
    {
      "name": "page view",
      "data": "{\"domain\":\"www.example.com\",\"title\":\"walkerOS documentation\",\"referrer\":\"https://www.walkeros.io/\",\"search\":\"?foo=bar\",\"hash\":\"#hash\",\"id\":\"/docs/\"}",
      "context": "{\"dev\":[\"test\",1]}",
      "globals": "{\"pagegroup\":\"docs\"}",
      "custom": "{\"completely\":\"random\"}",
      "user": "{\"id\":\"us3r\",\"device\":\"c00k13\",\"session\":\"s3ss10n\"}",
      "nested": "[{\"entity\":\"child\",\"data\":{\"is\":\"subordinated\"}}]",
      "consent": "{\"functional\":true}",
      "id": "ca81a74de5cb7d41",
      "trigger": "load",
      "entity": "page",
      "action": "view",
      "timestamp": "2023-11-14 22:13:20.123",
      "timing": 3.14,
      "source": "{\"count\":1,\"trace\":\"0a1b2c3d4e5f60718293a4b5c6d7e8f9\",\"type\":\"collector\",\"schema\":\"4\"}"
    }
  ],
  "format": "JSONEachRow",
  "clickhouse_settings": {
    "async_insert": 0,
    "input_format_skip_unknown_fields": 0,
    "input_format_null_as_default": 0,
    "insert_deduplication_token": "62e0b549e2b05fa504b7cab450bffc0764a399ff4aee069b6f9d20ab8e732933"
  }
})

Purchase

An order event is inserted as one row. The nested items array is JSON-encoded into the data column, so a richer payload needs no table change. Promoting a hot key to its own column is a MATERIALIZED expression in the table DDL, not a destination setting.

Event
{
  "name": "order complete",
  "data": {
    "id": "ORD-500",
    "currency": "EUR",
    "total": 199.99,
    "items": [
      {
        "sku": "SKU-1",
        "quantity": 2
      },
      {
        "sku": "SKU-2",
        "quantity": 1
      }
    ]
  },
  "context": {
    "shopping": [
      "complete",
      0
    ]
  },
  "globals": {
    "pagegroup": "shop"
  },
  "custom": {
    "completely": "random"
  },
  "user": {
    "id": "us3r",
    "device": "c00k13",
    "session": "s3ss10n"
  },
  "nested": [
    {
      "entity": "product",
      "data": {
        "id": "ers",
        "name": "Everyday Ruck Snack",
        "color": "black",
        "size": "l",
        "price": 420
      },
      "context": {
        "shopping": [
          "complete",
          0
        ]
      },
      "nested": []
    },
    {
      "entity": "product",
      "data": {
        "id": "cc",
        "name": "Cool Cap",
        "size": "one size",
        "price": 42
      },
      "context": {
        "shopping": [
          "complete",
          0
        ]
      },
      "nested": []
    },
    {
      "entity": "gift",
      "data": {
        "name": "Surprise"
      },
      "context": {
        "shopping": [
          "complete",
          0
        ]
      },
      "nested": []
    }
  ],
  "consent": {
    "functional": true
  },
  "id": "ae7f9a5b61413c0d",
  "trigger": "load",
  "entity": "order",
  "action": "complete",
  "timestamp": 1700000000456,
  "timing": 3.14,
  "source": {
    "count": 1,
    "trace": "0a1b2c3d4e5f60718293a4b5c6d7e8f9",
    "type": "collector",
    "schema": "4"
  }
}
Out
insert({
  "table": "events",
  "values": [
    {
      "name": "order complete",
      "data": "{\"id\":\"ORD-500\",\"currency\":\"EUR\",\"total\":199.99,\"items\":[{\"sku\":\"SKU-1\",\"quantity\":2},{\"sku\":\"SKU-2\",\"quantity\":1}]}",
      "context": "{\"shopping\":[\"complete\",0]}",
      "globals": "{\"pagegroup\":\"shop\"}",
      "custom": "{\"completely\":\"random\"}",
      "user": "{\"id\":\"us3r\",\"device\":\"c00k13\",\"session\":\"s3ss10n\"}",
      "nested": "[{\"entity\":\"product\",\"data\":{\"id\":\"ers\",\"name\":\"Everyday Ruck Snack\",\"color\":\"black\",\"size\":\"l\",\"price\":420},\"context\":{\"shopping\":[\"complete\",0]},\"nested\":[]},{\"entity\":\"product\",\"data\":{\"id\":\"cc\",\"name\":\"Cool Cap\",\"size\":\"one size\",\"price\":42},\"context\":{\"shopping\":[\"complete\",0]},\"nested\":[]},{\"entity\":\"gift\",\"data\":{\"name\":\"Surprise\"},\"context\":{\"shopping\":[\"complete\",0]},\"nested\":[]}]",
      "consent": "{\"functional\":true}",
      "id": "ae7f9a5b61413c0d",
      "trigger": "load",
      "entity": "order",
      "action": "complete",
      "timestamp": "2023-11-14 22:13:20.456",
      "timing": 3.14,
      "source": "{\"count\":1,\"trace\":\"0a1b2c3d4e5f60718293a4b5c6d7e8f9\",\"type\":\"collector\",\"schema\":\"4\"}"
    }
  ],
  "format": "JSONEachRow",
  "clickhouse_settings": {
    "async_insert": 0,
    "input_format_skip_unknown_fields": 0,
    "input_format_null_as_default": 0,
    "insert_deduplication_token": "7a52c5ac7dbc7eed04fdeb2a533113bf77611b74a7598db171b8d1f7bbb0f7b7"
  }
})

Credentials​

config.credentials carries the { username, password } of the ClickHouse user the destination writes as. Back the password with a managed secret ($secret.NAME) instead of writing it into the flow config. It wins over a username or password left in settings.clickhouse.

Passthrough options​

settings.clickhouse reaches the @clickhouse/client constructor, and settings.clickhouseSettings travels with every insert, so both stay open for options this destination does not model itself. Two rules apply on top:

  • url and database always win over the same keys in settings.clickhouse. An insert names only the table, so a passthrough database that disagreed would move every row to a different database in silence.
  • If you enable async_insert through clickhouseSettings, keep wait_for_async_insert = 1. At 0 the server acknowledges before the data reaches storage, so a failure at the later flush has no way back to the caller, reaches no retry and no dead letter queue, and is lost without an error anywhere.

Event columns​

Every event becomes one row. Payload fields travel as JSON-encoded strings, so the table needs no schema knowledge of what a flow sends.

ColumnEvent fieldTable typeNotes
idevent.idStringThe only input to the deduplication token.
nameevent.nameLowCardinality(String)entity action, with the space.
entityevent.entityLowCardinality(String)
actionevent.actionLowCardinality(String)
triggerevent.triggerLowCardinality(String)
timestampevent.timestampDateTime64(3, 'UTC')Written as a YYYY-MM-DD hh:mm:ss.SSS UTC literal.
timingevent.timingFloat64
dataevent.dataStringJSON, or an empty string when absent.
contextevent.contextStringJSON, or an empty string when absent.
globalsevent.globalsStringJSON, or an empty string when absent.
customevent.customStringJSON, or an empty string when absent.
userevent.userStringJSON, or an empty string when absent.
nestedevent.nestedStringJSON, or an empty string when absent.
consentevent.consentStringJSON, or an empty string when absent.
sourceevent.sourceStringJSON, or an empty string when absent.

No column is nullable. Inserts run with input_format_null_as_default = 0, so an absent payload arrives as an empty string rather than as a null the server would quietly replace with a default.

The timestamp is a basic-format UTC string rather than a number or an ISO literal on purpose. date_time_input_format defaults changed in ClickHouse 26.5 and integer parsing changed in 26.8, so a number can be read as Unix seconds by one server and as the raw value at the column's precision by another, which is a factor of 1000 apart. The basic format reads the same on every version a partner may run.

JSONEachRow is fixed and not configurable. A positional format would silently misalign every value the first time a column is reordered in a table the destination does not own.

Empty payloads​

An absent payload, an empty object and an empty array all become an empty string. An object whose keys are all undefined is not empty by that test and encodes as {}. Both forms exist in the table, so a query filtering out empty payloads should account for '' and '{}'.

Create the table​

There is no setup lifecycle and no provisioning. The destination creates nothing, and an insert into a table that does not exist fails with 60 UNKNOWN_TABLE.

Create the table before sending events. Replace {database} and {table} with the values in settings, and read the notes below before running it. The ORDER BY in particular is worth deciding first.

CREATE TABLE IF NOT EXISTS {database}.{table} (
  id String CODEC(ZSTD(1)),
  name LowCardinality(String),
  entity LowCardinality(String),
  action LowCardinality(String),
  trigger LowCardinality(String),
  timestamp DateTime64(3, 'UTC') CODEC(DoubleDelta, ZSTD(1)),
  timing Float64,
  data String CODEC(ZSTD(3)),
  context String CODEC(ZSTD(3)),
  globals String CODEC(ZSTD(3)),
  custom String CODEC(ZSTD(3)),
  user String CODEC(ZSTD(3)),
  nested String CODEC(ZSTD(3)),
  consent String CODEC(ZSTD(3)),
  source String CODEC(ZSTD(3)),
  user_hash String MATERIALIZED JSONExtractString(user, 'hash') CODEC(ZSTD(1))
) ENGINE = MergeTree
PARTITION BY toYYYYMM(timestamp)
ORDER BY (entity, action, timestamp)
TTL toDateTime(timestamp) + INTERVAL 14 MONTH DELETE
-- non_replicated_deduplication_window is what makes a retried insert safe on
-- this engine. It defaults to 0, which deduplicates nothing, so without it a
-- retry after a server-side timeout writes the same rows twice.
SETTINGS ttl_only_drop_parts = 1, non_replicated_deduplication_window = 1000

On a ReplicatedMergeTree the equivalent window is on by default and that setting is unnecessary. On this non-replicated MergeTree it is not optional: see Retries and duplicates for what the window covers and what it does not.

The lines worth explaining:

  • non_replicated_deduplication_window = 1000. The number of recently inserted blocks whose hashes the table keeps, which is what the insert_deduplication_token is matched against. At 1000 a retry is covered as long as fewer than 1000 further blocks landed in between. At the recommended batch settings one runtime inserts 0.2 to 1 blocks per second, so that is roughly 17 to 83 minutes, and every further runtime writing the same table shortens it in proportion. The destination retries within seconds, so this leaves a wide margin. Raise it if you insert far more often.
  • PARTITION BY toYYYYMM(timestamp). A partition key should stay low-cardinality, usually under 100 to 1000 distinct values. Monthly partitioning keeps a 14 month retention at 14 partitions. Partitioning finer (by day, or by site) is the common cause of 252 TOO_MANY_PARTS, because ClickHouse never merges parts across partitions, so each partition carries its own merge backlog.
  • ttl_only_drop_parts = 1. Retention then drops whole parts once every row in them has expired, instead of running the mutations a partial TTL cleanup performs. The server default is 0.
  • MergeTree, not ReplacingMergeTree. This destination writes raw hits and deduplicates at the insert layer, so a replacing engine buys nothing. It costs extra merge work, and it invites OPTIMIZE ... FINAL, which rewrites whole partitions and does not guarantee uniqueness at read time anyway.
  • 14 MONTH is a placeholder. Set the retention your contract requires.

Promoting a hot key to a column​

Payload columns are JSON strings, so a frequently filtered key becomes a real column through a MATERIALIZED expression over the string. It costs the destination nothing and needs no flow change:

ALTER TABLE {database}.{table}
  ADD COLUMN site_id LowCardinality(String)
  MATERIALIZED JSONExtractString(globals, 'site_id');

Two constraints:

  • A materialized column can only read a value the row actually contains, so the flow has to land the site on the event first (in globals, context or data). No table expression can invent it.
  • A column added later is computed for new parts only. Existing rows read the column's default until the parts are rewritten.

The ORDER BY is the decision to get right​

(entity, action, timestamp) is the safe default for a table with no tenant column. It is almost certainly wrong for an estate of thousands of sites: with that many tenants the leading column should be the site identifier, so that a query for one site reads one narrow range of the primary index instead of scanning every site's rows for the matching entity.

Changing the sorting key later is not an ALTER. It means creating a second table with the new key and moving the data across with ATTACH PARTITION, with ingestion paused for the switch. It is worth deciding before the first insert, and it depends on the same site value the materialized column above needs.

Preconditions​

Confirm these before the first event flows. Each one is a decision this package deliberately does not make for you:

  1. You run the DDL. The destination creates nothing. An insert into a table that does not exist fails with 60 UNKNOWN_TABLE.
  2. The ORDER BY prefix. See above. The default is safe, not right for a large multi-site estate, and expensive to change later.
  3. Deduplication is enabled. Either the table is a ReplicatedMergeTree, or it sets non_replicated_deduplication_window above 0. The reference DDL sets it, so this is a precondition on your own DDL: if you write your own, carry the setting across. See Retries and duplicates. Without one of the two, a retry can duplicate rows.
  4. The retention the contract requires. The 14 month TTL in the DDL is a placeholder.
  5. A loud failure is what you want. A walkerOS field that has no column in your table fails the whole batch rather than being dropped in silence. That trade is described in Schema drift fails loudly, and it means a schema change on our side can stop ingestion until you add the column.

Batching at scale​

Batching belongs to the collector, through config.batch. The package default is batching off, matching every other walkerOS destination, so the numbers below are a flow-config recommendation rather than a package default:

"batch": { "size": 5000, "age": 5000, "wait": 1000 }

ClickHouse's own guidance is explicit: "We recommend inserting data in batches of at least 1,000 rows, and ideally between 10,000-100,000 rows", and "We recommend keeping the number of insert queries around one insert query per second".

At roughly 300 events per second the settings above flush about 1,500 rows every 5 seconds, and at several times that load the size cap flushes 5,000 rows sooner. That is between 0.2 and 1 inserts per second across the range, which sits inside ClickHouse's recommendation at both ends.

The three fields do different jobs:

  • wait is a debounce. Its timer resets on every push, so it alone would never flush under continuous load.
  • age is the hard ceiling on how long the first event of a window waits. It is what bounds latency.
  • size is the hard ceiling on rows per insert. It is what bounds insert size at peak.

Retries and duplicates​

This destination's retry is the only retry in the pipeline. The collector has none: a delivery that throws is dead-lettered, and the dead letter queue is an in-memory buffer (100 entries by default, oldest dropped on overflow) that nothing drains back into the pipeline. Anything that reaches it is gone.

So a retry has to be safe. Every attempt of one batch sends the same insert_deduplication_token, a digest of the batch's event ids computed once before the first attempt. A server that had already written the block when the answer stopped coming skips the repeat.

The token only works if the table deduplicates. That needs one of:

  • a ReplicatedMergeTree family engine, where insert deduplication is on, or
  • non_replicated_deduplication_window set above 0 on a non-replicated MergeTree.

non_replicated_deduplication_window defaults to 0. On a plain non-replicated MergeTree with the default, nothing deduplicates and a retry can duplicate rows. The reference DDL sets it to 1000 for exactly this reason. If your table is your own, confirm which of the two covers you before enabling retries in production.

Two further conditions:

  • One batch stays one block. The window counts blocks, so a batch split into several blocks by a server-side setting is no longer covered by its own token.
  • The token replaces the block-content check. The server skips a repeated token while reporting success, with no error, no dead letter entry and no counter anywhere. That is why the token is derived from event ids and nothing else: a token that could repeat across different events would be silent data loss.

Retries are bounded by maxRetries (one retry by default, so two attempts) and by the derived request timeout, so the last attempt finishes inside config.timeout. Only transient failures are retried: server-side timeouts, a saturated query or connection pool, a cancelled query, an overloaded server, a table held read-only, and every transport failure that never got an answer at all (ECONNRESET and friends). Everything else is terminal and goes to the dead letter queue on the first failure.

Schema drift fails loudly​

Two settings travel with every insert and are the reason a missing column is noisy:

  • input_format_skip_unknown_fields = 0. The server default drops a JSON key the table does not know. At 0 it raises an error instead.
  • input_format_null_as_default = 0. The server default rewrites an explicit null into the column's default. At 0 it raises an error instead.

A walkerOS field with no column in your table therefore fails the whole batch, loudly, attributed and counted, rather than disappearing and being found weeks later in reporting. The cost is real: ingestion for that batch stops instead of degrading. If you would rather have the silent drop, set clickhouseSettings.input_format_skip_unknown_fields to 1 in the flow config. We think the loud failure is the right default for a warehouse of record, and this is the place to disagree with us.

Too many parts​

252 TOO_MANY_PARTS fires when the active parts in one partition exceed parts_to_throw_insert (3000 by default, with inserts delayed from 1000). It means merging cannot keep up with how many separate inserts are arriving.

The fix is a larger batch: raise config.batch.size so the same rows arrive in fewer inserts. Check the partition key too, since partitioning finer than monthly multiplies the merge backlog.

Raising parts_to_throw_insert is not the fix. It hides the backlog until it is worse.

Custom rows​

A mapping rule's data object replaces the canonical row wholesale. The object a mapping produces travels verbatim, so a table that is not the reference DDL is configured in the flow rather than in the destination:

{
  "destinations": {
    "clickhouse": {
      "package": "@walkeros/server-destination-clickhouse",
      "config": {
        "settings": {
          "url": "https://clickhouse.example.com:8443",
          "database": "analytics",
          "table": "orders"
        },
        "mapping": {
          "order": {
            "complete": {
              "data": {
                "map": {
                  "event_id": "id",
                  "event_name": "name",
                  "order_id": "data.id",
                  "revenue": "data.total",
                  "currency": { "key": "data.currency", "value": "EUR" }
                }
              }
            }
          }
        }
      }
    }
  }
}

With the matching table:

CREATE TABLE IF NOT EXISTS analytics.orders (
  event_id String,
  event_name LowCardinality(String),
  order_id String,
  revenue Float64,
  currency LowCardinality(String),
  inserted_at DateTime64(3, 'UTC') DEFAULT now64(3)
) ENGINE = MergeTree
PARTITION BY toYYYYMM(inserted_at)
ORDER BY (order_id, inserted_at)
-- A custom table needs the deduplication window just as much as the reference
-- one: the retry is the destination's and does not know your schema.
SETTINGS non_replicated_deduplication_window = 1000

Three things to know about a mapped row:

  • Every key must exist as a column. input_format_skip_unknown_fields = 0 applies here too, so a typo in a map key fails the batch.
  • The timestamp conversion does not apply. A mapped row is not built by the destination, so mapping "timestamp": "timestamp" sends raw epoch milliseconds into whatever the column is, which is exactly the version-dependent parsing the canonical row avoids. Fill the column with a server-side DEFAULT now64(3) as above, or produce the YYYY-MM-DD hh:mm:ss.SSS literal in the flow.
  • A batch is one insert, and a row can be short. A map key whose value resolves to undefined is left out of the row, and mapping only some events sends two shapes in one block. A column absent from a row is filled with its default rather than raising an error, which is the one silent path left. Give every column a default you can recognize, and keep one shape per destination: route a second shape to a second destination pointing at its own table.

The event id still comes from the event, never from the mapped row, so deduplication keeps working for a row that carries no id column.

Timeouts and shutdown​

config.timeout is the collector's per-delivery race (10000 ms when unset). The client's request timeout is derived from it so that every attempt and every retry delay finishes before the collector stops waiting.

At shutdown the collector races each pending batch flush against a 5 second step timeout, then closes the destination. A deployment that wants its retries to survive shutdown should keep config.timeout clearly below 5000, for example at 4000, so that the last attempt finishes inside that race with time to spare. At exactly 5000 the derived attempts and retry delays use all but a millisecond or two of the race, so ordinary timer delay can push the last attempt past it, and closing the destination then cuts off the insert still in flight.

Next steps​

💡 Need implementation support?
elbwalker offers hands-on support: setup review, measurement planning, destination mapping, and live troubleshooting. Book a 2-hour session (€399)