Metadata-Version: 2.4
Name: tkati-node-el
Version: 0.8.1
Summary: Generic extract/load node: reads from a configurable input and writes to a configurable output
Requires-Python: >=3.13
Description-Content-Type: text/markdown
Requires-Dist: tkati-core==0.8.1
Requires-Dist: loguru>=0.7.0
Requires-Dist: pydantic-settings>=2.11.0

# tkati-node-el — generic extract/load node

Reads batches from a configurable input and writes them to a configurable output. Offsets are committed only after the batch is written and confirmed delivered (`produce_arrow` followed by a blocking `flush`), so delivery is at-least-once.

The loop is `tkati-core`'s `PipelinedNode`: it reads the next batch while producing the current one, and doesn't wait for a batch's delivery before moving on. Batches are still committed only once delivered, in read order (up to `[pipeline] max_in_flight` batches, default 4, may be waiting). On SIGTERM or SIGINT the node finishes the batch in hand, waits for everything in flight to be delivered, commits it and exits. A second signal forces an exit, leaving uncommitted batches to be read again on restart.

Input and output kinds are selected via the `type` field in each section — pick from whatever `tkati-core` supports. Every backend's settings split a **`connection`** tier (server-specific: how to reach the broker/database) from the resource tier (`topic` for Kafka, `table` for ClickHouse) and, where relevant, a tier local to this reader/writer instance (Kafka's `consumer` settings).

- **Input**: `"kafka"` (JSON or Arrow-batch messages from a Kafka/Redpanda topic)
- **Output**: `"kafka"` or `"clickhouse"` (native Arrow insert)
- **DLQ**: same `OutputSettings` shape as `output` — a DLQ can be Kafka or ClickHouse too

## Configuration

Settings are loaded from a TOML file. Set the `SETTINGS_FILE` environment variable to point to it (defaults to `settings.toml`).

```toml
[input]
type = "kafka"

[input.connection]
broker = "redpanda:29092"

[input.topic]
name   = "traffic_event"

[input.topic.schema]
uid        = "string"
time       = "timestamp[ms]"
traffic_in = "uint32"
# … other columns

[input.consumer]
group_id          = "node-el-group"
batch_size        = 1000
batch_timeout_sec = 10
auto_offset_reset = "latest"

[output]
type             = "clickhouse"
dlq_split_factor = 10

[output.connection]
host     = "clickhouse"
port     = 9000
user     = "default"
password = ""
secure   = false

[output.table]
database = "default"
name     = "traffic_event"

[dlq]
type = "kafka"

[dlq.connection]
broker = "redpanda:29092"

[dlq.topic]
name = "node-el-dlq"

# Optional. Prometheus endpoint, ON by default; these are the defaults.
[metrics]
enabled = true
port    = 8000
addr    = "0.0.0.0"
```

A Kafka output instead looks like:

```toml
[output]
type = "kafka"

[output.connection]
broker = "redpanda:29092"

[output.topic]
name   = "some-other-topic"
format = "json"        # or "arrow-batch"
key_column = "uid"     # optional
```

A ClickHouse DLQ instead looks like:

```toml
[dlq]
type = "clickhouse"

[dlq.connection]
host     = "clickhouse"
port     = 9000
user     = "default"
password = ""
secure   = false

[dlq.table]
database = "default"
name     = "traffic_event_dlq"
```

## Metrics and perf log

Every 10 seconds the node logs where its wall clock went, using `LoopStats`
from `tkati-core`:

```
perf over 10s: 157000 rows in, 157000 out (0 dropped), 157 iterations (0 input-starved)
perf loop: wait/input=0.21s (2%) producer/serialize=2.10s (21%) producer/enqueue=0.35s (4%) producer/deliver=0.00s (0%) wait/in-flight=1.12s (11%) commit=0.38s (4%)
perf read: consumer/poll=4.43s (44%) consumer/parse=0.48s (5%) wait/loop=5.20s (52%)
```

This node never drops rows, so `dropped` is always 0.

`perf loop:` is the node's loop, in the order its phases happen for a batch.
`perf read:` is its read-ahead thread, which reads the next batches while the
loop works on this one. The two lines run at the same time, so add up
percentages within a line, never across them.

* `consumer/poll`, `consumer/parse`: fetching message batches from the broker,
  and decoding them into an Arrow table.
* `producer/serialize`, `producer/enqueue`: encoding rows into the output's
  wire format, and handing them to librdkafka.
* `producer/deliver`: a ClickHouse output records its whole insert here,
  retries and DLQ fallback included. A Kafka output doesn't wait for acks per
  batch, so it reads 0; see `wait/in-flight`.
* `commit`: the offset commit.
* `wait/input`: time the loop waited for the next batch from its read-ahead
  thread. High means the node is input-bound.
* `wait/in-flight`: time `done()` waited because `[pipeline] max_in_flight`
  batches were still undelivered. High means the node is output-bound.
* `wait/loop` (read line): time the read-ahead thread waited for the loop to
  take what it had read. High, with low `wait/input` and `wait/in-flight`,
  means encoding the output is the bottleneck.

Percentages are of the interval, so they **do not sum to 100**; the remainder
of a line is time in none of its named phases. `input-starved` counts iterations that
drained the topic and waited out the batch timeout. A mostly-starved interval
shows `consumer/poll` near 100% and says nothing about whether the node can
keep up. `tkati-core`'s README covers the phases in more depth.

The `perf over` and `perf loop:` numbers are served as Prometheus metrics at
`:8000/metrics` (the read line's are log-only). To turn
that off, set `[metrics] enabled = false` or the env var
`METRICS__ENABLED=false`; `METRICS__PORT` moves it. See `tkati-core`'s README
for the metric names and the PromQL that reproduces the log line's
percentages.

## DLQ semantics

DLQ *fallback triggering* is currently only implemented for the `clickhouse` output kind — `KafkaProducer` has no retry/split logic of its own. The DLQ *sink* itself (where isolated bad rows end up) can be Kafka or ClickHouse, independent of the primary output. When a batch insert fails after all retries, the app switches to a recursive fallback to isolate the problematic rows:

1. The failing batch is split into `dlq_split_factor` equal sub-batches and each is retried independently.
2. If a sub-batch also fails it is split again — this repeats until individual rows are reached.
3. A single row that ClickHouse still rejects is written to the DLQ sink, preserving the full schema (Arrow IPC `arrow-batch` format for a Kafka DLQ).
4. After all rows are handled (inserted or DLQ'd), the input offset is committed and the app resumes normal large-batch processing.

`dlq_split_factor` is a setting on the `clickhouse` `[output]` block (see above), not on `[dlq]` — it describes how the primary output retries, independent of where the DLQ sink sends isolated rows. With `dlq_split_factor=10` and a 1 000-row batch this takes at most 3 recursive levels (1000 → 100 → 10 → 1).

**Delivery guarantee: at-least-once.** If the process crashes mid-recursion the uncommitted batch is re-read on restart and re-processed from the beginning, which may produce duplicate rows in the output and duplicate messages in the DLQ.
