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.
# 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 downFlink 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)
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 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.
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.
- 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.
- Base image
flink:1.18.1-scala_2.12-java11; PyFlink pinned to the same1.18.1, image builtlinux/amd64(PyFlink'spemjahas no ARM64 wheel, so amd64 avoids a from-source build on Apple Silicon).