Overview
Every ten seconds a machine reports CPU and memory. The pipeline scores that sample as it arrives, with two detectors, so each stage can be timed on its own.
Context
ActivityReporter runs on my own machines. Every ten seconds a machine reports CPU and memory, and that sample is scored against the last hour of the same series before it is stored.
Publishing the sample is the small part. The project is the path after that: process the reading as it arrives, apply the statistics, and keep each stage separate enough to see what it costs. It is also where I am learning Kafka, Flink, and ClickHouse on a real stream rather than an exercise.
Problem
I wanted one concrete question at the end of that learning. Once live samples are moving through the pipeline, where do time and load actually go, stage by stage? The answer is a comparison of the components under the same load, not a dashboard of the machines themselves.
Constraints
A few home machines, one sample every ten seconds. Kafka, Flink, and ClickHouse are larger than that volume needs. That is deliberate: the tools are the subject, and the split between stages is what makes a component-by-component measurement possible.
The measurements are not in yet. Nothing here is a production result.
D-01Score each sample in the stream
- Decision
Each sample is published to Kafka. A Flink job, keyed by machine and metric, keeps the last hour of that series and scores the new value as it arrives. A separate writer batches the raw samples and the scores into ClickHouse.
- Why
Scoring in the stream is the processing problem. Keeping publish, score, and store apart is also what lets me time and load each one on its own, which is the result the project is for.
- Trade-off
More processes than a script that writes straight to a database. At this volume those processes are there to be measured and learned.
D-02Two detectors on one window
- Decision
Every scored sample is checked with a modified z-score (MAD) and with an exponentially weighted moving average (EWMA), both on the last hour of its own series. MAD flags a value far from the median in either direction. EWMA flags a rise above a moving average. Each emits a baseline and a flag, anomalous or not, so the chart can draw the band beside the raw line. A floor on the deviation keeps a quiet series from flagging every tiny step.
- Why
Two different definitions of unusual, on the same data. A new algorithm is a new label on the evaluation, not a schema change. The score joins back to its sample by id.
- Trade-off
The hour of state lives in the job. A restart clears the window and rescores from the start of the topic. That is acceptable while the stored history is a day; checkpointing waits until the measurements need a stable window.
D-03Store the series beside the score
- Decision
ClickHouse keeps the raw sample and one evaluation row per algorithm. Grafana draws both. The tables collapse a replayed row into one, and the samples and scores expire after 24 hours.
- Why
The reading has to show the baseline and the flag, not only the moments something looked wrong. A columnar store fits that append-only series, and it is the store I wanted to learn.
- Trade-off
More rows than a table of anomalies alone. A replay can show duplicates until the table merges them.
D-04When a stage stops
- Decision
The ClickHouse writer commits its Kafka offset only after the insert succeeds. If it dies mid-batch, that batch is read again. The table keeps one row per sample, so the second write does not become a second history. If the Flink job restarts, it scores the topic again from the beginning, and the same rule applies.
- Why
A dead process should hold the data up, not lose it. I can stop one stage and still recover the samples.
- Trade-off
Until ClickHouse merges, a replayed batch can show twice. Queries that need a single row use FINAL. The Flink job does not checkpoint yet, so a restart forgets the last hour and rebuilds it.
Outcome
In progress. Samples are scored end to end and visible in Grafana, with the baseline and the anomaly flag beside the raw series. What is still open is the comparison the project is for: how each stage performs under the same load. Final results will be posted here.
What I'd do differently
Too early to revisit a decision. I'll write this once the stage-by-stage numbers are in.