Skip to content

Latest commit

 

History

14 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Zero Copy Reconciliation Engine

The Zero Copy Reconciliation Engine is a high performance C++ financial data pipeline built to reconcile massive trade ledgers as fast as the hardware allows.

The engine is designed for workloads such as comparing two 10 million row trade files and identifying missing, mismatched, or divergent records. Instead of relying on conventional C++ file streams, heap heavy string parsing, and mutex based synchronization, this project takes a lower level systems approach focused on memory locality, cache efficiency, and parallel execution.

In benchmark tests, the engine reconciled 20 million total records in under 3 seconds in the nominal case, a 3.9x speedup over the same reconciliation written the conventional way with std::ifstream, std::string, and std::unordered_map.

The benchmark harness is part of the build. Running ./engine prints the table in Performance Benchmarks and computes the speedup ratios itself, so nothing in this README is hand-calculated.

This project is an applied study in hardware sympathy: writing software that works with the CPU, memory hierarchy, operating system, and storage layer rather than against them.

Why This Project Exists

Financial reconciliation is often treated as a data engineering problem, but at large scale it quickly becomes a systems performance problem. I built this project to explore how high throughput data pipelines can process massive transaction volumes without becoming a bottleneck on standard I/O, excessive allocations, lock contention, or inefficient data layouts.

A naive reconciliation pipeline typically spends a large amount of time on:

  • copying data between kernel and user space,
  • allocating millions of temporary strings,
  • chasing pointers inside std::unordered_map,
  • locking shared data structures across threads,
  • and repeatedly missing CPU caches.

This engine avoids those costs wherever possible.

The core idea is simple:

Map the files directly into memory, parse them without copying, hash the records into cache-friendly tables, and verify the second ledger in parallel using read only lookups.

Performance Benchmarks

Every row below processes the same 20,000,000 records — two 10,000,000-row ledgers. Measured on WSL2 Ubuntu 22.04 (ext4 inside the VHDX, not a /mnt/* Windows mount), Intel i7-1255U, 12 hardware threads, g++ 11.4, -O3 -march=native.

Scenario Records Median Best Worst
Baseline A — stream scan, first field only, 1 thread 20,000,000 805 ms 799 ms 1,612 ms
Baseline B — stream parse, all 4 fields copied, 1 thread 20,000,000 1,873 ms 1,795 ms 2,604 ms
Baseline C — full naive reconciliation, unordered_map, 1 thread 20,000,000 10,807 ms 10,345 ms 15,938 ms
Engine — worst case, complete ledger divergence 20,000,000 3,174 ms 2,608 ms 3,681 ms
Engine — nominal T+0 settlement 20,000,000 2,753 ms 2,443 ms 4,314 ms

Median-run speedup vs Baseline C: 3.93x. Across the observed spread the ratio ranges from 2.40x (worst engine run against the best baseline run) to 6.52x (the reverse).

The Nominal T+0 benchmark includes the full pipeline:

  1. memory-mapping the input files,
  2. zero-copy CSV parsing,
  3. hashing and indexing the internal ledger,
  4. merging thread-local maps,
  5. scanning the external ledger,
  6. and verifying all records against the global index.

Reading the table honestly

Baseline C is the only fair comparison, and it is the only one the headline number uses. It runs the same algorithm as the engine — index the internal ledger, then verify the external one against it with the identical break rule — using std::ifstream, std::string, and a reserved std::unordered_map, single threaded. Both implementations independently report the same 5 breaks on the nominal dataset, which is the strongest available evidence that the engine's speed does not come from skipping work.

The engine is slower than Baselines A and B, and that is expected. Those two only parse and count; they build no index and perform no lookups, so they finish in 805 ms and 1,873 ms while the engine spends its time doing work they never do. They are published because they bound the problem: A is the floor for touching 970 MB of CSV at all, and the A-to-B gap (about 1,070 ms) is the raw cost of the ingestion copy path that string_view removes.

Part of the 3.93x is thread count, not zero-copy. The engine uses 12 threads; every baseline is single threaded. This harness does not attempt to separate the two contributions, and the number should not be quoted as though it were purely an I/O or allocation result.

Roughly 2.6 GB of fixed allocation is inside the timed region. Each run allocates 12 thread-local FlatMaps at 1 << 20 slots plus a global map at 1 << 24 slots, at 88 bytes per slot, and that zero-fill is counted against the engine. It is fixed cost that does not scale with input size — at 200,000 rows it alone made the engine slower than the baselines. Leaving it inside the measurement is deliberately conservative.

Methodology

The harness is explicit about how it measures:

  • One discarded warmup, then N timed runs, applied identically to every row. This is not optional. On an identical workload the first engine run measured 3.4x slower than the third (5,665 → 2,531 → 1,650 ms) purely from cold page faults against that 2.6 GB of table. Without a warmup, whichever scenario runs first looks bad and whichever runs last looks good.
  • Headline ratios use the median run, not the best. The engine swings roughly 65% run to run on this hardware while Baseline C is stable to within about 5%, so best-of-N lets a single lucky engine rep set the published number. The median is the statistic that holds still.
  • The page cache is warmed before every timed run, so no row pays for cold disk reads the others avoid.
  • Record counts are measured, not assumed. The engine reports global_map.count() (unique internal trade IDs indexed) plus external rows scanned. The 20,000,000 in the table is a number the program counted.
  • Baseline allocations are pinned with a compiler barrier. At -O3, GCC can prove that only the length of each parsed field is ever consumed, fold substr(0, n).size() to n, and delete the string allocations outright — which would leave the baselines timing getline alone. do_not_optimize() prevents that.
  • Filesystem matters more than anything else here. These numbers are only valid on native ext4, where this machine scans the ledgers at roughly 1.2 GB/s. The same reads across a /mnt/* Windows mount under WSL run more than an order of magnitude slower, and because that penalty falls hardest on the stream-based baselines, it silently inflates any speedup measured there. Do not benchmark this project from /mnt/c or /mnt/d.

Reproducing

python3 scripts/generate_data.py                   # ~1m45s, writes 2 x 485 MB
cp ./data/trades_int.csv ./data/trades_ext_some.csv
sed -i '1,5s/^\([^,]*\),[^,]*,/\1,999.99,/' ./data/trades_ext_some.csv   # inject 5 price breaks
make clean && make
./engine 7                                          # arg = timed reps per scenario, default 3

Needs about 8 GB of RAM: 1.5 GB of ledger data in page cache, ~2.6 GB of engine tables, and ~2.5 GB for Baseline C's unordered_map. Expect roughly 3 minutes at 7 reps.

Core Architecture

1. Memory-Mapped File Input

Instead of reading CSV files through std::ifstream or repeated read() calls, the engine uses POSIX mmap().

This maps the entire file directly into the process address space. The operating system still manages paging and caching, but the application no longer needs to manually copy file contents into intermediate buffers. This reduces overhead from:

  • repeated system calls,
  • user-space buffering,
  • redundant memory copies,
  • and stream abstraction costs.

The result is a file input path that behaves much closer to direct memory access.

2. Zero Copy Parsing with std::string_view

Traditional CSV parsers often allocate a new std::string for every field in every row. At 10 million rows, that can mean tens of millions of heap allocations. This engine avoids that!

Each parsed field is represented as a std::string_view, which is only a lightweight pointer-length pair referencing the original memory-mapped file. No field text is copied during ingestion.

For example, a trade ID is not copied into a new string. Instead, the engine stores a view into the mapped file:

std::string_view trade_id;

This makes parsing much cheaper and keeps the memory footprint predictable.

3. Thread Local Parsing and Lock Free Verification

The reconciliation process is split into phases.

Phase 1: Build Local Indexes

The internal ledger is divided into chunks. Each worker thread parses its own chunk and inserts records into a private thread-local hash map. Since each thread owns its own map, there is no shared write contention and no need for mutexes during ingestion.

Phase 2: Merge into a Global Index

After all threads finish parsing, the local maps are merged into one global index. This happens in a controlled phase rather than during parsing, which keeps the hot ingestion path free from locks.

Phase 3: Parallel Verification

The external ledger is then parsed in parallel. Each thread performs read-only lookups against the global index. Since the global index is no longer being modified, lookups can happen without locks.

This design gives the engine a clean map reduce style architecture:

Internal Ledger
      ↓
Parallel Thread Local Indexing
      ↓
Global Merge
      ↓
External Ledger Verification
      ↓
Reconciliation Results

4. Cache Friendly Flat Hash Map

A standard std::unordered_map is convenient, but it is not ideal for this workload. Most implementations use node based storage, where each entry may live at a separate heap location. This causes poor cache locality and frequent pointer chasing.

This engine uses a custom open addressed flat hash map instead. Records are stored in contiguous memory, and collisions are resolved using linear probing. This greatly improves cache behavior because nearby entries are likely to be loaded together.

Key design choices include:

  • Open addressing instead of node based chaining,
  • linear probing for predictable memory access,
  • power of two capacity for fast indexing,
  • bitwise masking instead of modulo division,
  • and MurmurHash3 for strong 64-bit hash distribution.

Instead of index = hash % capacity; the engine uses index = hash & (capacity - 1);. This replaces an expensive division operation with a fast bitwise operation.

Setup

Requirements

You will need:

  • Linux on a native Ext4 filesystem — under WSL this means the distro's own filesystem, never /mnt/c or /mnt/d, which will distort every number in the benchmark,
  • GCC or G++ with C++20 support (tested on g++ 11.4),
  • Python 3 for generating test data,
  • about 8 GB of free RAM,
  • and a CPU that benefits from -march=native.

This project is designed for Linux because it uses POSIX APIs such as mmap().

Running the Project

Generate the Internal Trade Ledger using

python3 scripts/generate_data.py

This creates two large synthetic internal and external trade files under the data/ directory.

Next, for the nominal reconciliation benchmark, copy the internal ledger:

cp ./data/trades_int.csv ./data/trades_ext_some.csv

To simulate reconciliation breaks, change a few prices or quantities in trades_ext_some.csv. Doing it with sed rather than by hand keeps the expected break count a known quantity:

sed -i '1,5s/^\([^,]*\),[^,]*,/\1,999.99,/' ./data/trades_ext_some.csv

Both the engine and Baseline C should then report exactly 5 breaks. If they disagree, one of them is wrong.

Now, we can compile the Engine:

make clean && make

The Makefile builds the project using C++20 and native CPU optimizations.

Finally, run the benchmark:

./engine        # 3 timed reps per scenario (default)
./engine 7      # 7 reps, less noise, ~3 minutes

The executable runs all three baselines and both engine scenarios, then prints the summary table and the speedup ratios.

Current Capabilities

The engine currently supports:

  • parsing large CSV trade ledgers
  • memory mapped file access
  • zero copy field extraction
  • thread-local indexing
  • global hash map merging
  • parallel external ledger verification
  • mismatch detection
  • a self-contained benchmark harness with three single-threaded baselines, warmup, repeat runs, and median-based speedup reporting
  • cross-validation of results against an independent naive implementation

Future Optimization Roadmap

  • Size the hash tables to the input. This is the largest measured win available, and the benchmark is what surfaced it. Every run allocates a fixed ~2.6 GB of table (12 thread-local maps at 1 << 20 slots plus a global map at 1 << 24, 88 bytes per slot) regardless of how much data is actually being reconciled. At 200,000 rows that zero-fill alone made the engine slower than a single-threaded getline loop. Sizing the maps from the file length, or allocating them once and reusing them across runs, would remove a fixed cost that currently sits inside every measurement. The Slot struct itself is also worth shrinking: at 88 bytes, one slot spans more than a cache line.

  • Transparent Huge Pages: The engine can also experiment with huge page advice (madvise(ptr, length, MADV_HUGEPAGE)). Using 2 MB pages instead of standard 4 KB pages can reduce TLB misses when scanning very large files and hash tables.

  • SIMD CSV Scanning: The next major optimization is SIMD based delimiter scanning. Instead of checking one character at a time, AVX2 or AVX-512 instructions could compare 32 or 64 bytes in parallel when searching for commas and newlines. This would speed up CSV tokenization, especially on wide records.

  • NUMA Aware Thread Pinning: On multi socket systems, memory access can become slower when a thread reads memory owned by another NUMA node. A future version could use pthread_setaffinity_np to pin worker threads to specific physical cores and reduce cross socket memory traffic. This would make performance more stable on larger servers.

Key Takeaways

This project demonstrates that high performance data processing is not only about using more threads.

The largest gains come from removing unnecessary work:

  • avoid copying bytes that do not need to move,
  • avoid allocating objects that do not need to exist,
  • avoid locks on hot paths,
  • avoid pointer heavy data structures,
  • and keep memory access predictable for the CPU.

The result is a reconciliation engine that treats financial ledger processing as a systems level performance problem and solves it using control over memory, concurrency, and cache friendly data layout (i.e. flat hash table layout).

About

High performance C++ financial data pipeline built to reconcile massive trade ledgers

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages