Igor Kołodziej

Scala · FS2 · Redpanda · PostgreSQL · MongoDB

Payment Event Pipeline

A Scala 3 backend that can replay payment events locally or consume them from Redpanda without forking the processing logic. Events are validated, enriched with customer data, evaluated against stateful eligibility and risk rules, and persisted as replay-safe projections.

Architecture diagram showing JSONL replay, paced replay and Redpanda converging on one EventSource, followed by parsing, validation, enrichment, historical context, eligibility and risk evaluation, and MongoDB persistence. PostgreSQL supplies customer profiles and MongoDB feeds historical state back into context building.
File replay and broker-backed ingestion share the same processing path. Historical state enters decision logic through explicit interfaces.

One processing model for replay and streaming

The pipeline has a single input contract: an EventSource produces a stream of records with a source position and raw payload. File replay, paced replay and Redpanda are adapters behind that interface. Once an event enters the application layer, its source no longer matters.

That keeps the core flow simple: parse, validate, enrich, build context, decide, persist. FS2 provides back-pressure and streaming composition, while Cats Effect keeps I/O explicit. The same pipeline can therefore be exercised deterministically from a local JSONL file and then switched to broker-backed ingestion through configuration rather than a second code path.

Stateful decisions without burying state in the rules

Several decisions depend on recent customer behaviour: transaction velocity, repeated failures, device history, unusual amounts or shifts in payment-method usage. The rules themselves do not query MongoDB. A RiskContextProvider builds the relevant history first and passes a compact context into the decision engine.

This separates data access from decision logic and makes the rules straightforward to test. Eligibility is evaluated independently from risk: an inactive account, insufficient balance or currency mismatch is a business-rule decline, while behavioural signals contribute to a separate risk assessment. The final decision combines those two layers explicitly rather than collapsing every failure into a generic score.

Replay-safe persistence

Replaying historical input is only useful if it does not duplicate the persisted result. Processed events are written with an upsert keyed by eventId. Alerts use (eventId, alertType). Eligibility violations use (eventId, violationType). Matching unique MongoDB indexes are created by the application on startup.

The stored event history is also the input to later stateful decisions. Queries are scoped by customer and time window, with the current event excluded from its own history. Stable persistence codes for domain enums keep the stored representation independent from Scala case names, so refactoring application code does not silently change database semantics.

The less visible engineering work

A large part of making the service reliable was tightening the contracts around the happy path. Parsing and validation return explicit domain errors. Validation can accumulate multiple problems in one event. Customer and event IDs use opaque types. Money is kept in explicit currencies rather than converted implicitly. Startup configuration is parsed into typed settings with invalid values rejected early.

The PostgreSQL adapter uses Doobie behind a customer-lookup port, MongoDB access is isolated behind storage and history interfaces, and the local environment is reproducible with Docker Compose. The repository also has unit tests around the pipeline and decision logic, formatting checks, CI and coverage reporting. These details matter because they make the architecture enforceable in code rather than just visible in a diagram.

From this design to production

The current implementation already has the pieces that are hardest to retrofit later: explicit boundaries, deterministic replay, idempotent writes and testable decision logic. The next production steps are mostly about delivery guarantees and operations.

  • tie Kafka offset progression to durable processing, or introduce an inbox/outbox pattern for stronger delivery guarantees.
  • add structured logs, metrics and tracing around ingestion latency, decision outcomes and persistence failures.
  • run PostgreSQL, MongoDB and Redpanda integration tests automatically with disposable containers.
  • version the external event schema and define partitioning and ordering guarantees explicitly.

Those changes would strengthen the operational model without requiring the domain and processing layers to be rewritten, which was one of the main architectural goals of the project.