Skip to content

Ingest Plugin

Use @chkit/plugin-ingest to load application API data into tables defined in your chkit schema.

Start with the API sync quickstart for a complete first sync. Read the API sync guides for reader, storage, and checkpoint choices, or install the authoring skill for a coding agent.

  • Runs finite pulls from application APIs into chkit-managed tables.
  • Records independent run histories in an append-only ClickHouse journal, supporting overlapping executions.
  • Records progress after writes succeed. A retry can reread rows from the last committed checkpoint.
  • Retries source requests with backoff, Retry-After, and failure classification.
  • Selects streams by exact tags so any external scheduler (cron, CI, Kubernetes) can drive it.

Create and change destination tables through schema migrations before running ingestion.

Set entry in place of schema globs. chkit imports the entry module once and collects its exported tables and pipelines. Export a pipeline to activate it; remove its export to deactivate it.

Configure a direct clickhouse connection for ingestion. Host-provided executors, including the ObsessionDB workbench executor, do not provide the required JSON row encoding and per-insert deduplication settings.

clickhouse.config.ts
import { defineConfig } from '@chkit/core'
import { ingest } from '@chkit/plugin-ingest'
export default defineConfig({
entry: './src/chkit.ts',
plugins: [ingest()],
clickhouse: { url: process.env.CLICKHOUSE_URL ?? '' },
})

The target needs CREATE, SELECT, and INSERT permissions for the journal and destination tables. ChKit stores checkpoints entirely in ClickHouse journal evidence.

Every run records its starting checkpoint and local event sequence. Event identity includes the run, so overlapping executions and stale reads cannot reuse another run’s identities. Complete snapshot validation selects acknowledged progress from valid histories. Malformed tails are diagnosed and excluded; normal overlapping runs need no repair. A stale valid snapshot may replay work, so destination keys and versioning must reconcile repeats.

See Scheduling and recovery for execution and recovery behavior.

A stream is a destination table plus a read generator. Fetching and mapping are ordinary code inside read; each yielded chunk carries rows already shaped for the table.

src/chkit.ts
import { table } from '@chkit/core'
import { HttpError, definePipeline, defineStream, ingestionColumns, paginate, timestampWindow } from '@chkit/plugin-ingest'
export const tickets = table({
database: 'crm',
name: 'tickets',
columns: [
{ name: 'id', type: 'String' },
{ name: 'updated_at', type: "DateTime64(3, 'UTC')" },
{ name: 'raw', type: 'String' },
...ingestionColumns,
],
engine: 'ReplacingMergeTree(updated_at)',
primaryKey: ['id'],
orderBy: ['id'],
})
const ticketStream = defineStream({
id: 'helpdesk.tickets',
destination: tickets,
tags: ['schedule:1h'],
incremental: timestampWindow({
// Re-read one hour of overlap; ReplacingMergeTree reconciles repeats.
start: new Date(0),
overlapMs: 3_600_000,
}),
async *read(context) {
const pages = paginate({
context,
label: 'GET /tickets',
fetchPage: async (cursor: string | undefined, signal) => {
const url = new URL('https://api.example.com/tickets')
// This example assumes the provider supports both range parameters.
url.searchParams.set('updated_since', context.selection.from.toISOString())
url.searchParams.set('updated_before', context.selection.to.toISOString())
if (cursor) url.searchParams.set('cursor', cursor)
const response = await fetch(url, { signal, headers: { Authorization: `Bearer ${process.env.HELPDESK_TOKEN}` } })
if (!response.ok) throw await HttpError.fromResponse(response)
const body = await response.json() as {
data: Array<{ id: string; updated_at: string }>
next_cursor: string | null
}
return { items: body.data, next: body.next_cursor ?? undefined }
},
})
for await (const page of pages) {
yield { rows: page.items.map((item) => ({ id: item.id, updated_at: item.updated_at, raw: JSON.stringify(item) })) }
}
},
})
export const helpdesk = definePipeline({ id: 'helpdesk', streams: [ticketStream], maxFetches: 4 })

Add ingestionColumns to custom destination tables. The loader fills _chkit_batch_id and _chkit_run_id; ClickHouse sets _chkit_ingested_at at the physical insert.

paginate returns complete pages shaped as { items, next, metadata? }, including empty and terminal pages. Use page.items for rows. Pass checkpoint candidates from metadata explicitly to a yielded chunk’s state when using cursorState; metadata itself is not persisted. See Readers and pagination.

Batch identity decides whether a retry is deduplicated. By default it includes a content hash of the rows, which prefers a possible duplicate over suppressing rows that changed between attempts; a field like synced_at: new Date() therefore defeats retry deduplication. When a chunk covers a stable source interval, declare it with id (for example yield { rows, id: \page:${cursor}` }`): the chunk then becomes its own write unit and its identity ignores row content.

Retain raw objects when the query shape may change and storage is affordable. Map in the reader when the destination schema is established. rawTable stores the returned provider object in a native JSON column beside a stable id; use rawRows to load a page and a SQL view to expose typed fields:

import { view } from '@chkit/core'
import { HttpError, defineStream, paginate, rawRows, rawTable } from '@chkit/plugin-ingest'
export const rawTickets = rawTable({ database: 'crm', name: 'tickets_raw' })
export const tickets = view({
database: 'crm',
name: 'tickets',
as: `SELECT
id AS ticket_id,
raw.subject::String AS subject,
raw.requester.email::String AS requester_email,
arrayMap(t -> t.name::String, raw.tags[]) AS tags,
parseDateTime64BestEffortOrNull(raw.updated_at::String, 3, 'UTC') AS updated_at
FROM crm.tickets_raw FINAL`,
})
const ticketStream = defineStream({
id: 'helpdesk.tickets',
destination: rawTickets,
async *read(context) {
const pages = paginate({
context,
fetchPage: async (cursor: string | undefined, signal) => {
const url = new URL('https://api.example.com/tickets')
if (cursor) url.searchParams.set('cursor', cursor)
const response = await fetch(url, {
signal,
headers: { Authorization: `Bearer ${process.env.HELPDESK_TOKEN}` },
})
if (!response.ok) throw await HttpError.fromResponse(response)
const body = await response.json() as {
data: Array<{ id: string; [key: string]: unknown }>
next_cursor: string | null
}
return { items: body.data, next: body.next_cursor ?? undefined }
},
})
for await (const page of pages) yield { rows: rawRows(page.items, (ticket) => ticket.id) }
},
})

The raw table uses ReplacingMergeTree to retain the latest ingested version per id. Query with FINAL to resolve repeats before background merges finish. Changes to an ordinary view can use retained fields without re-fetching the source. A materialized view needs a backfill to update stored results; fields you did not retain require a source re-fetch. See Destinations and transformations.

StrategyUse whenBookmark advances
none (full sync)The source is small or has no change filterNever; every run reads everything
timestampWindow({ start, overlapMs })The API filters by an updated-since timestampTo the run cutoff, after the whole window loaded
cursorState({ id, version, parse })The provider owns the state: compound cursor, change token, page positionWhenever a yielded chunk carries state and its rows have been saved

timestampWindow starts its first sync at start. Later runs begin at the committed watermark minus overlapMs (zero by default). Explicit backfill bounds take precedence. For custom lower bounds, use timestampWindow({ from: ({ watermark, cutoff }) => ... }) instead. Both forms retain the same checkpoint format and only advance after the entire window is saved.

Provider clients can accept the exported FetchContext type, containing attempt and signal. ReadContext extends it with the stream selection and checkpoint, so readers can pass their context directly without coupling clients to checkpoint generics.

With cursorState, state on a chunk must be the complete state that is safe to resume from once every row up to that chunk is saved. Omit it when you cannot make that claim; the run then restarts from the previous checkpoint after a failure.

A checkpoint records its strategy id and version. Changing either makes the next run fail rather than reinterpret old state.

Pipeline retry settings are defaults; a stream can override individual settings. Source attempts and reader restarts use p-retry for exponential backoff, jitter, retry counts, and maxRetryTime. The supported options are retries, factor, minTimeout, maxTimeout, randomize, maxRetryTime, shouldRetry, and shouldConsumeRetry.

The policy callbacks receive p-retry’s attemptNumber, retriesLeft, and retriesConsumed, plus a FetchFailure preserving the original cause and normalized classification. Returning false from shouldConsumeRetry follows p-retry’s behavior: it skips consuming a retry and skips its backoff. A provider’s Retry-After still applies, including for those unconsumed retries. Retry waits release fetch and load capacity for other streams.

Once a source attempt exhausts its policy, it fails the stream; reader recovery does not multiply that retry budget. Opaque reader failures restart from the latest committed checkpoint.

Terminal window
chkit ingest list # streams in the loaded graph
chkit ingest run # run every stream
chkit ingest run --tag schedule:1h # exact tag match; repeat --tag for AND
chkit ingest run --tag stream:helpdesk.tickets
chkit ingest status # committed checkpoint per stream
chkit ingest doctor --tag stream:helpdesk.tickets
chkit ingest repair --tag stream:helpdesk.tickets
chkit ingest run --backfill jan --from 2026-01-01 --to 2026-02-01

Every stream also carries the derived tags pipeline:<id> and stream:<id>. schedule:<cadence> is a convention only: chkit never interprets it. A --tag filter that matches nothing fails before any work runs.

A backfill uses its own checkpoint namespace, so it never moves the scheduled bookmark. Reusing its ID reuses that state, but resumption depends on the strategy: explicit timestamp bounds take precedence over the watermark and reread that range. The bundled fullSync() and cursorState() strategies do not interpret date bounds; custom provider strategies may honor or reject them. See the source’s documentation and Backfill source data.

doctor reads and validates evidence without changing tables. repair previews recovery by default. After reviewing the plan, pass its fingerprint to repair --apply <fingerprint> with exactly one selected stream. Add --backfill <id> to either command to inspect or repair that historical namespace. With writers stopped, repair preserves evidence and materializes a separately validated journal. Set ingest({ journalTable: '<replacementTable>' }) to the returned table before restarting writers; activation is explicit. Uncertain work may be replayed. See the command reference for details.

chkit check verifies that every stream destination carries the ingestion metadata columns. See chkit ingest for the full flag reference, exit codes, and JSON output.

Ingestion is at-least-once. The loader inserts batches with stable insert_deduplication_token values. ClickHouse suppresses retries after a lost acknowledgement when the table’s engine, settings, and deduplication window support it. Pick a destination engine that reconciles repeats for your data, for example ReplacingMergeTree keyed by the provider ID.

Successful syncs start a new batch identity cycle, recorded by the existing journal. Failed or interrupted syncs retain their cycle for replay. This also applies to full syncs, which have no incremental bookmark.

The duration budget bounds journal operations as well as readers. Shutdown gives unfinished readers or writes up to five seconds to settle. Terminal journal appends are bounded by the duration budget while the run is live. After cancellation or budget exhaustion, each one gets a separate five-second limit. Interrupted writes never count as successful ingestion.

Overlapping processes record independent histories. Scheduler concurrency controls limit resource use; a stale valid read may replay acknowledged work. ChKit validates journal facts before using their checkpoint.

OptionDefaultDescription
journalTable_chkit_ingestion_journalAppend-only checkpoint journal in the configured database
maxDurationSeconds3600Execution budget. Exhausting it ends the run as incomplete and keeps committed progress
prefetchBatches1Mapped batches buffered between fetching and loading