← Projects

StockPulse — Stock Anomaly Detection & AI Analysis Pipeline

  • Mar 2026
  • Solo project · Design to operations
  • Kafka · MSA
  • Kubernetes · ArgoCD
  • TimescaleDB

Overview

When a price suddenly jumps during trading hours, the first question is “why?”. StockPulse watches 1-minute bars for 68 US and 40 Korean tickers, and when a sharp move appears it first decides whether the move belongs to one stock, a whole sector or the entire market. It then gathers related news, has an LLM explain the cause in Korean and English, and delivers it to a dashboard and Slack.

From ingestion to alerts, eight services are connected by seven Kafka topics (two of them DLQs) in an event-driven design, deployed and operated on an on-prem Kubernetes cluster with ArgoCD GitOps. It is a solo project I designed, built and ran myself.

Architecture

Scroll and the diagram on the right follows the text on the left. Hover over a node to see its specs.

At a glance

StockPulse takes in external data (prices, news, an LLM), processes it with services that talk over Kafka inside an on-prem Kubernetes cluster, and delivers results to the browser and Slack. Deployment runs through the GitOps flow at the bottom.

services ×8topics ×7Kafka ×3

1How prices come in

For 68 US tickers, stock-collector pulls 1-minute bars from yfinance every 60 seconds during NYSE hours and publishes them to stock.raw.us. When the market closes it sleeps until the next open.

stock.raw.usNYSE calendar

For 40 Korean tickers, kis-bridge receives live trades over the Korea Investment WebSocket, rolls them into 1-minute bars and publishes them to stock.raw.kr. It renews tokens before they expire and reconnects with a 5→60 s backoff.

stock.raw.krKIS H0STCNT0

2Detecting and classifying moves

anomaly-detector reads both topics and keeps only bars from the last 5 minutes whose change or Z-score crosses a threshold. Korea gets a higher threshold because of its ±30% daily limit.

% OR Z-scorelast 5 min

ETFs act as sector thermometers. Three or more sector ETFs moving the same way means MARKET; the sector’s ETF or its peers moving together means SECTOR; otherwise INDIVIDUAL. Results go out on anomaly.detected.

INDIVIDUALSECTORMARKET

3Finding the cause

news-fetcher consumes anomaly.detected, gathers up to four articles per language from the last three days via NewsAPI, Google News and Naver RSS using sector keywords and the ticker, and publishes news.fetched.

anomaly.detectednews.fetched

ai-analyzer writes Korean and English analyses with the Groq LLM, shifting the prompt’s focus to company, sector or macro issues by class. On a 429 it waits as instructed and retries; if it still fails, the message goes to news.fetched.dlq.

llama-3.3-70bcircuit breaker 5×/60snews.fetched.dlq

4Reaching people

api consumes analysis.completed, stores it in TimescaleDB and pushes it straight to open dashboards over the /ws/live WebSocket. Browsers come in through the nginx Ingress.

analysis.completed/ws/live

notifier filters the same results by sector and event type and sends them to Slack. After three failed retries the message moves to analysis.completed.dlq.

3 retriesanalysis.completed.dlq

Restart behavior differs by stage. Detection, news fetching and notifications use earliest to reprocess anything missed; LLM analysis and api use latest so a backlog doesn’t trigger a burst of LLM costs or re-show the same signals.

earliestlatest

5A model that retrains daily

The ml-trainer CronJobs predict next-day direction for 25 Korean tickers. Each day they grade yesterday’s calls and retrain if 30-day accuracy drops below 52% or the model is older than 7 days. Once a week they tune with Optuna and validate walk-forward.

LGB 0.4 · XGB 0.35 · Cat 0.25Purged CV

6From git push to deploy

A push to main triggers GitHub Actions on a self-hosted runner, which maps changed paths to services, builds only what changed and pushes stock/<svc>:<sha7> to Harbor.

path → service mapimages ×9

Actions commits the new image tags to the manifests ([skip ci] prevents loops), and ArgoCD checks the repo every 60 seconds and makes the cluster match Git with selfHeal and prune. Images are pulled from Harbor, and secrets live in Git encrypted with Sealed Secrets.

selfHeal · pruneServerSideApply

7Keeping watch — observability

Prometheus scrapes api and anomaly-detector metrics plus Kubernetes and node metrics every 15 seconds, and one Grafana dashboard shows Kafka lag, API p95 latency, detections and ML accuracy.

anomaly_detected_totalml_prediction_accuracy
At a glance1 / 14
External dataOn-prem Kubernetes · namespace stockUsersCI/CD · GitOpsKAFKA ×3yfinanceUS 1-min bars · 60 sKIS OpenAPIKR live tradesNewsAPI · RSSEN · KR newsGroq LLMllama-3.3-70bSlackAnomaly alertsBrowserLive dashboardstock-collectorUS · market hourskis-bridgeticks → 1-min barsnews-fetcherup to 4 per languageai-analyzerKO/EN cause analysisnotifierfilters · retriesanomaly-detector% · Z · ETF classesapiFastAPI · WebSocketTimescaleDBcompressed chunksfrontendNext.js 14Ingressnginx · MetalLBml-trainerCronJob ×3PrometheusGrafana · 15 s scrapeSealed Secretsdecrypts secretsGitHubmain branchActionsself-hosted runnerHarborprivate registryArgoCDselfHeal · prune

Try it

A demo where recorded prices flow through the Kafka pipeline and sharp moves are detected and classified

Coming in the next update.

Operations log

ArgoCD stuck in OutOfSync forever

Symptom
Some resources stayed OutOfSync after every deploy, and large manifests failed to apply at all.
Cause
Fields the cluster fills in — StatefulSet volumeClaimTemplates, CronJob status — never matched Git, and large resources exceeded the 262 KB limit of the last-applied annotation that client-side apply leaves behind.
Fix
Excluded those fields with ignoreDifferences and switched to ServerSideApply=true so resources apply without the annotation.

With selfHeal and prune enabled, ArgoCD keeps syncing as long as a difference remains. Differences caused by cluster-populated fields never go away no matter how often you sync, so the dashboard was always OutOfSync and real drift was hard to spot. I declared which fields to ignore, capped retries at five (from 5 s up to 3 min), and moved deletions last with PruneLast.

Lesson In GitOps, "equal to Git" has to leave out fields the cluster owns. Design what you compare and how you apply together.

Three brokers eyeing a small node’s memory

Symptom
Three brokers’ heaps added up to nearly 6 Gi, enough to risk OOMKills on small worker nodes.
Cause
Default JVM settings for the stateful workloads were sized for bigger machines, and startup and shutdown times weren’t accounted for.
Fix
Capped broker and ZooKeeper heaps at -Xmx256m, gave the startupProbe up to 630 s to wait for __consumer_offsets loading, and raised terminationGracePeriodSeconds to 90 s for a clean shutdown.

Kafka runs three brokers (RF 3, min ISR 2). Alongside memory, I tuned startup and shutdown: a starting broker needs time to read __consumer_offsets, and if a probe declares failure first, the pod restarts and repeats the same work. The startupProbe now allows enough time, and shutdown waits for the controlled shutdown to finish.

Lesson For stateful workloads, the time it takes to come up and go down is part of resource design.

Minute-bar batches outgrew Kafka’s message limit

Symptom
Sending 1-minute bars for 68 US tickers as one message hit the message size limit.
Cause
The collector packed the whole lookup window (5 days by default locally) into a single batch message.
Fix
The Kubernetes deployment uses a 1-day window, and the message limit is set to 10 MB.

Every 60 seconds the collector sends all tickers’ minute bars at once under the key batch. Locally it sent five days, but in the cluster it sends only one day to stay under the message size limit.

Lesson "One batch = one message" runs into limits as data grows. Decide how to split messages that can get big.

Design decisions

A different consumer offset strategy per stage

Context
Each pipeline stage wants different behavior after a restart: some must reprocess what they missed, others must not.
Alternatives considered
  • earliest for every consumer
  • latest for every consumer
Why
Detection, news fetching and notifications use earliest, so a restart reprocesses missed messages. LLM analysis and api use latest, so a restart doesn’t trigger a burst of LLM costs on the backlog or re-show the same signals on the dashboard.
Trade-off
latest stages skip messages from while they were down. Signals that missed analysis can be re-run with the dashboard’s "reanalyze past" action.

Classifying signals with ETFs as sector thermometers

Context
When one stock jumps, whether it is company news or a sector-wide move changes both the cause and the news worth reading.
Alternatives considered
  • Judge each stock alone and leave cause-finding to the LLM
Why
Each sector has ETFs as thermometers: three or more sector ETFs moving the same way means MARKET; the sector’s ETF or peers moving together means SECTOR; otherwise INDIVIDUAL. The class also shifts the LLM prompt’s focus to company, sector or macro issues.
Trade-off
Classification quality depends on ETF coverage and thresholds, which need tuning per market (higher in Korea because of the ±30% daily limit).

Wait, cut off and quarantine external API failures in a DLQ

Context
External APIs like the LLM and Slack can slow down or return 429 at any time. One message’s failure must not block the whole pipeline.
Alternatives considered
  • Retry a failed message until it succeeds
  • Drop failed messages
Why
ai-analyzer reads the wait time in a 429 response and retries after it, and after five consecutive failures it stops calling for 60 seconds (circuit breaker). Messages that still fail go to news.fetched.dlq with the original topic, error and timestamp. notifier likewise sends to analysis.completed.dlq after three Slack retries.
Trade-off
There is no tool yet to replay messages from the DLQs.

A fallback path that runs without Kafka

Context
If everything depends on Kafka, the collectors and the detector all being up, losing part of the infrastructure stops the whole service — and checking features locally needs the full stack.
Alternatives considered
  • Refuse to start without Kafka
Why
When KAFKA_BOOTSTRAP_SERVERS is empty, api runs a daily-bar pipeline (detect → news → analyze → store) itself every hour with APScheduler. Both paths use the same detection and analysis modules (core/), so the logic doesn’t drift.
Trade-off
The fallback works on daily bars, so it isn’t real time. In production the streaming path is the default.

GitOps with image-tag commits and ArgoCD

Context
Building nine services’ images and rolling them out by hand makes it hard to know what is running at which version.
Alternatives considered
  • Deploy straight from CI with kubectl apply
Why
CI builds only changed services, pushes them to Harbor and commits the new image tags (commit SHAs) to the manifests, with [skip ci] preventing loops. ArgoCD sees the commit and makes the cluster match Git with selfHeal and prune, so Git history answers what is deployed.
Trade-off
Tag commits pile up in the history, and ArgoCD’s own constraints (annotation limits, cluster-populated fields) needed separate handling (see the operations log).

Secrets in Git, too — Sealed Secrets

Context
With every manifest in Git for GitOps, the question is where secrets go.
Alternatives considered
  • Create Secrets by hand with kubectl
  • An external secret store (e.g. Vault)
Why
A manually triggered workflow builds a .env from GitHub Secrets, encrypts six SealedSecrets with kubeseal and commits them. Only the in-cluster controller can decrypt them, so they are safe in Git. The controller itself is managed as an ArgoCD Application.
Trade-off
Losing the controller key means recreating every SealedSecret, so key backup has to be part of operations.

Infrastructure & cost

git pushmainActionsself-hosted runnerHarborstock/<svc>:<sha>Tag commit[skip ci]ArgoCDselfHeal · pruneClusternamespace stock

On-prem Kubernetes

  • A single stock namespace runs 12 Deployments, 3 StatefulSets (Kafka ×3, ZooKeeper, TimescaleDB), 2 DaemonSets, 6 CronJobs and 3 Jobs.
  • Traffic enters through one MetalLB (L2) address and an nginx Ingress: /api, /ws and /auth go to FastAPI, everything else to Next.js. WebSocket timeouts are raised to one hour.
  • Storage combines static NFS PVs (Retain) with nfs-subdir-external-provisioner, and heavy workloads (Kafka, the database, the API, ML) are steered to node-tier=heavy nodes.
  • Images live in a private Harbor, and a DaemonSet distributes its self-signed certificate to each node’s containerd.

From git push to deploy

GitHub Actions on a self-hosted runner maps changed paths to services and builds only the services that changed. It pushes stock/<svc>:<sha7> to Harbor and commits the new tags to the manifests; ArgoCD checks that commit every 60 seconds and applies it to the cluster. Secrets are stored in Git encrypted with Sealed Secrets.

Retention and cost

  • TimescaleDB hypertables compress chunks older than a day, and drop prices after 30 days and analyses after 90 by whole chunks — far lighter than row-by-row DELETEs.
  • A daily CronJob handles chunk cleanup, VACUUM ANALYZE and a size report; weekly CronJobs watch NFS usage (warning at 80%) and rewrite the Redis AOF.
  • Kafka keeps at most 7 days or 15 GiB.

My contribution

This section is coming soon.

Limits & retrospective

Known limitations, and what I would change next time.

  • The API is effectively a single instance. WebSocket connections are tracked in process memory, so running several API replicas means some browsers miss updates. The HPA allows up to three, but it has to stay at one until broadcasts move to Redis pub/sub.
  • CI has no tests. The pipeline builds only what changed, but there is no test stage before the build, so regressions are caught by hand.
  • Partition counts are not set. Topics are auto-created, so partition counts are implicit and adding consumers helps less than it should.
  • Gaps in observability. ai-analyzer is a scrape target without a metrics server, kafka-exporter isn’t deployed so the Kafka lag panel may be empty, and Redis is deployed but unused.
  • CronJob time zone. The jobs set timeZone: Asia/Seoul but their schedules were written in UTC, so training runs nine hours earlier than intended.