Skip to content

Latest commit

 

History

2 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

streaming-starter

A minimal Apache Flink (PyFlink) streaming job running on a real Flink cluster (JobManager + TaskManager) via the official Flink Docker image. It's the warm-start for a real-time market-data pipeline — the primitives here (source → keyBy → stateful aggregation → sink) are what the full project extends.

Run it

# 1. Build the image + start the Flink cluster (JobManager + TaskManager)
docker compose up -d --build

# 2. Submit the PyFlink job to the cluster
docker compose exec jobmanager flink run -py /opt/job/job.py

# 3. See the output — the print sink goes to the TaskManager's stdout,
#    which this console-mode cluster surfaces in its container logs.
#    PyFlink prints records as Flink Rows (+I[...]), so grep loosely:
docker compose logs taskmanager | grep -iE 'aapl|msft|goog'

# 4. Stop everything
docker compose down

Flink Web UI: http://localhost:8081 — watch the job run, see the job graph, and view Task Managers → (the TM) → Stdout for the sink output directly.

Expected output (running per-symbol totals as events arrive):

('AAPL', 1)
('MSFT', 1)
('AAPL', 2)
('AAPL', 3)
('MSFT', 2)
('GOOG', 1)
('AAPL', 4)
('MSFT', 3)

What it does

synthetic trade events  →  keyBy(symbol)  →  running sum(qty)  →  print

job.py is ~30 lines and deliberately minimal — source, keying, stateful aggregation, sink. The real work in Phase 3 is a swap-in (live feed + event-time windowing), not a rebuild.

The one thing this taught me (worth knowing)

v0: a processing-time tumbling window.

The first version used a processing-time tumbling window and it emitted nothing. Reason: the bounded source produces all events in milliseconds and the job finishes long before a 5-second wall-clock timer fires, so the window never triggers. Real windowing on bounded/finite data needs event time + watermarks, which fire deterministically on end-of-input.

v1: a watermark is a periodic signal, not a per-event calculation

I predicted the out-of-order trade ("MSFT", 9, 3s) would be dropped as late: it arrives after a trade at 26s, far past the 5s bound. It was counted instead (MSFT [0,10) = 17, not 8).

Why: Flink emits watermarks on a wall-clock timer (pipeline.auto-watermark-interval, default 200 ms), not after every record. The bounded source pushed all 12 trades through in a few milliseconds, before the first tick, so the watermark never advanced mid-stream.

Where this is going (roadmap)

  • v0 — hello world: bounded source, keyed running sum, on a real Flink cluster. (this)
  • v1 — event-time windowing: assign timestamps + watermarks, tumbling event-time windows per symbol (fires correctly on bounded and live data).
  • v2 — live source: replace the bounded list with a continuous feed (SEC EDGAR API or a generated market-data stream).
  • v3 — serve + harden: expose results via a small API; add an architecture diagram, cost notes, and a "what I'd change at 100k events/sec" section.

Notes

  • Base image flink:1.18.1-scala_2.12-java11; PyFlink pinned to the same 1.18.1, image built linux/amd64 (PyFlink's pemja has no ARM64 wheel, so amd64 avoids a from-source build on Apple Silicon).

About

A minimal Apache Flink (PyFlink) streaming job on a Dockerized Flink cluster, built as the warm-start for a real-time market-data pipeline.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Contributors

Languages