ML Infra/Platform
Scoring 39 million rows in 36 minutes: Optimizing batch inference at Intuit Credit Karma
August 10, 2025
Sarthak Gupta
Software & ML Engineer
The first version of my batch scorer produced the wrong scores. It was a few hundred lines of Apache Beam that read a BigQuery table, loaded a pickled model on every Dataflow worker, and called predict_proba one row at a time. On a small sample it looked fine. Then I diffed its output against the production scorer, and the numbers didn't agree. The first time I pointed it at the full 39-million-row table, it ran out of memory. And it was built on a model format the org's security preferences steered away from.
That V1 was supposed to replace the scorer behind a big part of Credit Karma. A lot of what members see in the app and in their inbox starts as a model score. On a regular cadence, the platform scores tens of millions of rows against models that produce "algorithmic facts": short, personalized insights that feed in-app stories and email. One of those scoring flows read about 39 million rows and roughly 5 TB of features. It ran for 6 hours and 44 minutes and cost $603 per run.
That was tolerable for a couple hundred models, but it stopped being tolerable when Data Science asked to onboard ten times as many.
Over one summer on the ML Platform team, we turned that broken V1 into a production batch scoring component built on Apache Beam, Google Cloud Dataflow, and PMML. The same job now finishes in 36 minutes for $89.54: 92% faster and 85% cheaper, with output that matches the legacy scorer down to the 17th decimal place. That puts projected savings at more than $400K a year and leaves room for 10x more models without 10x the bill.
Getting there took about 3 months: weeks of making two pipelines agree on the same numbers, a detour through Java version archaeology, forty JVMs quietly eating 686 GB of RAM, and dozens of benchmark runs that each tested one change. This post walks through all of it, including the dead ends.
A 7-hour job standing between Data Science and 10x more models
Batch scoring is the least glamorous part of an ML system. A model gets trained once, but it gets scored over and over against fresh features for every member, and the outputs land in BigQuery tables that downstream products read. When scoring is slow, three things happen:
- Facts go stale. Algorithmic facts are only as fresh as the last scoring run. A seven-hour job can't be refreshed often, so the stories and emails built on it lag behind the data.
- New models get rationed. Data Science had a backlog of models it wanted to put into production, and every one meant another multi-hour, several-hundred-dollar job. More models meant more revenue, but on the legacy scorer they also meant compute and cost growing at the same rate.
- The scorer couldn't be reused. Scoring lived inside a monolithic pipeline that was hard to change, and upcoming work like the Tax ML launch needed a batch scorer it could call directly.
Runtime and cost both grew linearly with the number of models, so the 10x ask would have turned a one-model refresh into a multi-day, roughly $6,000 job.
So the mission was simple to state: let the team score models in a fraction of the time and cost. We set four hard constraints before writing any code, and every decision in the rest of this post traces back to one of them:
- Identical scores. A faster scorer that disagrees with production is not acceptable, i.e. having parity with the legacy scorer was a non-negotiable. Not even a difference of 0.0001 was allowed.
- A fraction of the cost, not just the time. Throwing bigger machines generally does not help, but I did run a benchmark for this in-case it did the trick (hint: it didn't).
- A reusable component. Packaged, containerized, and launchable by any team, especially easy to use for Data Scientists inside their notebooks, rather than another step bolted onto the monolith.
- GCP-native. Built on the Beam/Dataflow and BigQuery stack our ML platform already ran on, so there was no new infrastructure to operate.
Carving a batch scorer out of a monolith
Two systems matter throughout this post.
Vega is the legacy ML pipeline framework. In Vega, batch scoring is one stage inside PipelineProcessor, a large, multi-file processing engine kicked off from Airflow. An Airflow preprocessing query materializes features into a BigQuery table, with many features packed into a single JSON column. PipelineProcessor reads that table, applies Vega's defaults and imputation, and scores each row through a custom wrapper around the JPMML evaluator.
VegaX is the newer, modular ML SDK the platform is migrating to. It had components for training and explanation, but no modern batch scorer.
The VegaX batch scorer is a standalone Apache Beam pipeline. For readers who haven't used Beam: a pipeline is a graph of transforms over distributed collections (PCollections). The per-element work happens in DoFns, and a runner, here Google Cloud Dataflow, fans that work out across a pool of worker VMs. We package the pipeline as a Dataflow Flex Template, a Docker image plus a launch spec, so anyone can start a scoring job by passing an input table, a model path, and an output table.
Getting the new pipeline to produce exactly the same output as the old one turned out to be the hardest part of the project.
Why I refused to publish until two pipelines agreed to the last digit
My rule from day one was that I wouldn't publish anything until VegaX produced the same scores as Vega for the same input. If VegaX ran faster but scored differently, it would effectively be a new model that nobody had validated, so speed comparisons were meaningless until the scores matched.
So before tuning anything, I set parity up as a gate. Vega became the ground truth and VegaX the candidate, in a shadow setup. Both read the same source table with the same model version. Their outputs were joined on numericId, and calc_prob_class1 (Vega) was diffed exactly against predicted_proba_class_1 (VegaX), every digit and every row. A change could move on to benchmarking only if the diff came back empty. If anything differed, the change was blocked and I went back to the inputs.
The first time I ran the gate against V1, the diff wasn't close.
Five suspects for one wrong score
When two ML pipelines disagree, the list of suspects is long, so I worked through it in order of how cheap each one was to rule out:
- Row order. I grouped and sorted by
numericIdbefore comparing, and it made no difference. - Model version. Both pipelines were confirmed to load the same artifact.
- Library versions. XGBoost and
sklearn2pmmldifferences can shift scores slightly. They don't produce the gaps I was seeing. - Scoring logic. Comparing Vega's
pmml_scorer.pywith VegaX's much simplerbatch_scorer.pyturned up real logic differences, but none that explained everything. - The features themselves.
The fifth one was the answer. It only showed up after a Zoom session with Pallav, Staff MLE, where we decided to stop comparing outputs and compare inputs instead: narrow the test to a single member and dump every feature value each pipeline handed to the model.
| Feature | Vega | VegaX (first pass) | |
|---|---|---|---|
| TRDLN4_util062 | null | 0 | mismatch |
| INQR207_cnt18m | 71 | null | mismatch |
| INQR207_cnt03m | 3 | 3 | match |
| feature count | from model metadata | from table columns | mismatch |
| calc_prob_class1 (after fixes) | 0.13861947205318462 | 0.13861947205318462 | identical |
Some features matched exactly, which proved the underlying data was identical. Others were blank in one pipeline and populated in the other. That pattern pointed at how each pipeline read and prepared the data.
Four plumbing bugs disguised as a bad model
Unwinding that diff took about a week and turned up four separate problems, each of which looked like "the model is wrong":
The feature list came from the wrong place. VegaX was building its feature vector from whatever columns the input table had. Vega used the list stored in the model's own metadata. The counts didn't match, so VegaX was effectively feeding the model a differently shaped input. I switched VegaX to read features from the model's meta_data.
Features were locked inside a JSON column. Many features live in a packed JSON column in BigQuery. Vega unpacked them, but VegaX's first pass didn't, so those features reached the model as nulls and got imputed. Fixing the unpacking alone closed roughly 65–70% of the gap.
Some features didn't exist in the table at all. A handful of fields the model expected, such as Relative_date and recsys_itemid, were created by Vega's Airflow preprocessing query, not stored in the source table. When I first ran VegaX at full scale on Dataflow, it failed on exactly these fields. I rebuilt the input table from the same query Airflow used, so both pipelines read identical inputs.
None is not NaN. This was the subtle one. scikit-learn's imputers expect missing values as np.nan. Python's None either crashed scoring or slipped past the imputer. BigQuery can't store NaN in the first place, so Vega encodes missing values as a sentinel (-1) and maps them back before scoring. VegaX's ImputeAndCastType step needed to reproduce that round trip exactly: sentinel to null to NaN, cast to the type the model expects.
The common thread in all four fixes was making the model artifact the source of truth. The feature list now comes from the model's metadata, and missing values are encoded the way the model was trained to read them. Later, the output schema was derived from the model as well. After those changes, the two pipelines stopped drifting apart.
After those four fixes, the same member scored 0.13861947205318462 in both pipelines, every digit identical. I later repeated the exercise after moving to PMML, this time on a sample of 500 numericIds, and the gate stayed in place for every change after that.
Why I threw out a 10x speedup on day one
With parity in place, I ran VegaX at full scale on Dataflow for the first time:
- Input: 38,921,433 rows
- Output: 38,921,433 rows
- Runtime: 36 minutes
Compared with Vega's 6 hours 44 minutes, that looked like a 10x win on day one. It wasn't a fair comparison, though. To get past V1's OutOfMemoryError, I had moved to n2d-highmem-64 workers (64 vCPUs, 512 GB of RAM each), twice the size of Vega's n2d-standard-32. Another job died for a much sillier reason: a debug beam.io.WriteToText I'd forgotten to comment out. Part of any speedup on bigger hardware comes from the hardware itself, so I couldn't credit all of it to the pipeline.
So before optimizing anything, I built the boring tool that made the rest of the project possible: a benchmark tracker. Every full-scale run went into one sheet with its job name, worker type, worker count, batch size, JVM settings, runtime, estimated cost, and a note on what changed. I couldn't get direct access to GCP billing exports, so cost was estimated from worker-hours and machine pricing. That stayed consistent across runs, which was what mattered. Dataflow runtimes also vary by a few minutes between identical runs, so any "winning" configuration got rerun before I believed it.
The first question for the tracker was whether hardware alone was the answer. Pallav and I agreed to get a clean baseline first, so I gave Vega the same n2d-highmem-64 machines:
- Vega on
n2d-standard-32: 6h 44m, $602.74 - Vega on
n2d-highmem-64: 3h 22m, $717.18
Doubling the machine halved the runtime and raised the cost by about 19%. Brute force made the job faster, but it made the cost problem worse. Whatever made VegaX fast had to come from the pipeline itself.
Trading pickle for PMML, and Python for a JVM
V1 loaded the model as a Python pickle and called predict_proba. That worked, but pickles have a well-known problem: loading one can execute arbitrary code. The organization preferred PMML, a declarative XML format for models that a scoring engine interprets rather than executes. Vega already scored PMML through its own wrapper, so to match Vega and meet the security bar, VegaX had to score PMML too.
The standard PMML engine is JPMML, a Java library. Its Python binding, jpmml_evaluator, starts a JVM inside the Python process through pyjnius, a JNI bridge, and every evaluation call crosses from Python into Java and back. That boundary comes up again in almost every section below.
Debugging a JVM inside a Python process inside a container
For consistency, I started on the same JPMML evaluator version Vega pinned. It immediately failed on the models with InvalidElementException. The models declared more input fields than the evaluator would accept. Newer versions offer a lax mode that tolerates this. The older one didn't.
Upgrading to 0.15.1 fixed that, but it required Java 11. Java 11 removed JAXB from the JDK, and the newer evaluator still needed it, so the next failure was a NoClassDefFoundError. Once I added a JAXB runtime, the Docker image failed because pyjnius couldn't find libjvm.so, which was fixed by setting JAVA_HOME and JVM_PATH explicitly in the image.
What made this painful was that everything passed locally. My workstation still had a cached Java 8 toolchain, so local runs used a different Java than the container and the Dataflow workers. Every bug appeared only after a Docker build and a cloud submit, and each loop took tens of minutes. A company-wide GCP incident in the middle of it took BigQuery and CI down for most of a day, which I spent debugging the JDK setup locally.
| Evaluator + runtime | Where it ran | Result |
|---|---|---|
| jpmml-evaluator 0.2.x + Java 8 | Vega's pinned version; also what my laptop had cached | Fails. No lax mode, so models with more inputs than fields throw InvalidElementException. |
| jpmml-evaluator 0.15.1 + Java 11 | What the Docker image and Dataflow workers actually ran | Fails. Lax mode works, but JAXB left the JDK in 11, so it throws NoClassDefFoundError. |
| 0.15.1 + Java 11 + JAXB runtime + explicit JVM paths | Rebuilt image with JAVA_HOME and JVM_PATH set and libjvm.so on the path | Works locally, in Docker and on Dataflow. |
From then on I tested inside the same container image that production used. Once local, Docker, and Dataflow all ran the same Java, the first full-scale PMML run finished in 39 minutes. That was correct, but slower than the pickle version, and it was where the real optimization started.
Benchmarking a data pipeline like an ML experiment
From here I ran performance work as a series of controlled experiments: change one variable, run the full ~39M-row job, put it through the parity gate, record it in the tracker, and keep or revert the change. The runs fell into four themes: brute force (above), batching, work around the model, and JVM memory. Figure 7 shows the runs that mattered.
Amortizing the Python-to-Java toll with BatchElements
JPMML's Python evaluateAll() still evaluates records one at a time on the Java side. Calling it once per element from a Beam DoFn meant paying the full overhead on every row: convert a Python dict into a Java map, cross the JNI boundary, run the evaluator, and convert the result back. With 39 million rows, that fixed cost dominated.
Beam's BatchElements transform groups elements into lists before they reach a DoFn. That let me build one DataFrame per batch and make one evaluator call per few hundred rows instead of one per row. I also moved model loading out of a Beam side input and into DoFn.setup(), so each worker parses the PMML file and starts its JVM once and caches the evaluator for the rest of the job:
class GenerateScores(beam.DoFn):
def __init__(self, pmml_model_path):
self.pmml_model_path = pmml_model_path
self.evaluator = None
def setup(self):
# Once per worker process, not once per element:
# parse the PMML, start the JVM, cache the evaluator.
self.evaluator = create_model_evaluator(self.pmml_model_path)
def process(self, batch):
features = pd.DataFrame(batch)
yield from score_batch(self.evaluator, features) # one JNI round trip per batch
scores = (
rows
| "Batch" >> beam.BatchElements(min_batch_size=300, max_batch_size=300)
| "Score" >> beam.ParDo(GenerateScores(model_path))
)
At a batch size of 1,000, the job ran in 39m 13s for $129.45, down from $603. Going higher didn't help. At 1,500 the job got slower and more expensive (40m 55s, $135.85): bigger batches hold more memory at once, each batch takes longer, and the runner has fewer, chunkier units of work to spread across workers. Removing logging inside the hot loop didn't move the needle either.
Smaller batches were far worse. In an early probe on a deliberately small heap, Dataflow projected the scoring stage alone would take about 90 days at a batch size of 1, and still about 6 days at 100. Those runs mixed two effects, small batches and a starved heap, so I don't treat them as clean measurements. They were still enough to show that most of the cost was per-call overhead, which batching could spread across many rows.
The best batch size also depended on the memory settings, which comes up again in the JVM section below.
Profiling found the bottleneck outside the model
Profiling Dataflow's per-step wall time pointed at a transform called MapVegaDefaultsToNull, which took longer than scoring itself. It converts Vega's sentinel default values back to nulls. It did this for every feature column in the table, even though the model only consumes a subset, its forward features. Restricting it to forward features, which are now read from the model rather than hardcoded, took most of that step's cost away.
Two smaller changes followed the same idea of not carrying work through the pipeline that the model never uses:
- Extract the JSON features once. The pipeline now unpacks the packed JSON column at the top of the pipeline instead of re-parsing it downstream.
- Carry fewer columns. It only keeps the columns the output table needs. Wide rows cost serialization time at every step in a distributed pipeline.
Together with batching, these brought the PMML pipeline from 38 minutes to 28 minutes on the highmem workers. That was the fastest runtime I recorded. I ended up shipping a slightly slower configuration, for reasons covered in the next section.
Forty JVMs that each thought they owned a 512 GB machine
The 28-minute pipeline had a problem. It passed on a single row and then died with an out-of-memory error at full scale. My first suspects were the change that had moved the model off a side input and the imputation step, so I bisected: comment out imputation, rerun, compare memory graphs with the last good job. Hardik, an engineer on the team, also pointed out that I was hitting resource exhaustion in the shared GCP project by running Vega and VegaX benchmarks side by side, so I stopped running jobs in parallel. Neither change fixed the OOM.
The problem turned out to be the JVM's default heap sizing. Every Dataflow worker starts its own JVM through pyjnius to host the evaluator, so 40 workers means 40 JVMs. Without explicit limits, a JVM sizes its maximum heap from the host's physical memory, which on a 512 GB machine is a very large number. It doesn't allocate that memory up front. It grows into it as garbage accumulates, defers collection while there's headroom, and shares the machine with a Python process that also needs memory. Eventually the two together exceed the box.
The fix was to give every worker an explicit memory budget instead of letting each JVM size itself. It lives in the Docker image, and the line that matters is the last one:
ENV JAVA_HOME=/usr/lib/jvm/jdk-11
ENV JVM_PATH=${JAVA_HOME}/lib/server/libjvm.so
# Bound every worker's JVM instead of letting it size itself from the host.
ENV JAVA_TOOL_OPTIONS="-Xms256m -Xmx2560m"
Getting to that number took several iterations. The first ceiling, 12 GB, cut the job's aggregate RAM use from 686 GB to 191 GB. That told me most of that memory was heap the JVM had claimed but scoring never used. From there I kept tightening: 6 GB, then 3 GB, then 2.5 GB, checking runtime and cost at each step.
Capping the heap changed the economics. The job no longer needed highmem machines at all:
- 6 GB heap,
n2d-highmem-64, batch 1,000: 32m 42s, $104.47 - 2.5 GB heap,
n2d-standard-64, batch 300: 36m 00s, $89.54
The second run is about three minutes slower and 14% cheaper, and I shipped it. A smaller heap does better with smaller batches, since each batch's DataFrame and its Java-side copies have to fit comfortably under the ceiling. At 300 rows per call, the per-call overhead was already amortized. I also stopped short of the lowest heap that would pass. A production job that runs out of memory at 3 a.m. costs more than a few gigabytes of headroom, so I left a buffer.
From pipeline to platform component
The last stretch of the project turned the pipeline into a component other teams could use.
A Flex Template anyone can launch. The Docker image bundles Beam on Python 3.10, JDK 11 with the JAXB runtime, and the JPMML evaluator, and the template takes the input table, model path, config, and output table as parameters. The parity gate became an automated test: it launches a real Flex Template job through the Dataflow API and checks the output table against Vega's.
A CI bug hiding in commit messages. Template launches started failing at random with missing-image and manifest errors. The cause turned out to be in CI: certain commits, notably the "format code" commits from running black and isort, caused the pipeline to skip the base image build. The template then pointed at an image tag that had never been pushed. Tracking it down took two days, and the fix itself was small.
An output schema the model defines. Vega writes a probability_class_* column for every class the model predicts. A ProbaColumnAugmenter now reads the model's output fields and derives the BigQuery schema from them, so the same component handles binary and multiclass models without code changes, and the output can't drift from what the model produces.
A long review. The PR collected dozens of review comments from senior engineers, and it got much better for it. Forward features are fetched dynamically instead of hardcoded. Project IDs and column schemas moved into constants and config. Redundant helpers like a homegrown flatten_json were removed. The Dockerfile dropped a vendored JDK tarball. The legacy apply_sklearn2pmml_fixes path was removed, since PMML scoring made it dead code. Pallav and I briefly considered merging in two halves, core scoring first and multiclass columns second, but we chose to keep iterating on one PR so the component never landed half-finished.
The batch scoring component is now merged into the VegaX main branch. It's headed for the shared core ML libraries next, and it will be the scorer behind the Tax ML launch.
Five rules I'd hand the next team migrating a scorer
Gate every change on parity. About a third of the calendar went to parity before a single optimization landed. It felt slow at the time, but it's the reason I could trust every benchmark number afterward.
Benchmark like a scientist. One variable per run, every run in the tracker, winners rerun before they're believed. The tracker turned "it feels faster" into a table I could defend in review.
Measure whether bigger machines actually help. Doubling Vega's hardware made it faster and more expensive. The large wins came from removing work: per-call overhead, unused columns, a transform running on features the model never reads.
Know where your runtimes meet. The most expensive things in this pipeline were boundaries: Python to Java, Python's None to scikit-learn's NaN, BigQuery's lack of NaN to the model's expectations. Most of the bugs and most of the speedups were at those seams.
Test where you ship. A cached Java 8 on one laptop cost three days. If production runs in a container, test in that container.
Next: taking Python out of the hot path
Native Java scoring. Every batch still crosses the Python–JVM boundary through pyjnius. I prototyped a Beam cross-language (external) transform that hands rows straight to a Java PMML evaluator, and a Beam version upgrade to support it has been cleared. The open problems are batching across the language boundary, which requires Beam's schema-aware Row objects instead of plain dicts, and a remaining score mismatch caused by how Python-side features map onto the inputs the Java evaluator expects. Until that mismatch is fixed, the parity gate keeps it out of production.
GPUs for batch inference. The platform has GPU capacity that isn't heavily used. The team has started exploring NVIDIA's Triton Inference Server for offline scoring, alongside Credit Karma's existing CPU-based inference server, Prophecy.
Wider rollout. Next is moving the component into the shared core ML libraries and running it for the Tax ML launch. After that comes the backlog of models that were waiting on a faster scorer.
The bottom line: $15.49 to $2.30 per million rows scored
Back to the four hard constraints from the start of this post. Every number below comes from the same ~39M-row, ~5 TB input on 40 workers.
| Metric | Vega (legacy) | VegaX (shipped) | Change |
|---|---|---|---|
| Wall-clock time per run | 6h 44m | 36m | −92% |
| Compute cost per run | $602.74 | $89.54 | −85% |
| Cost per 1M rows scored | $15.49 | $2.30 | −85% |
| Throughput | ~1,600 rows/s | ~18,000 rows/s | 11.2x |
| Workers | 40 × n2d-standard-32 | 40 × n2d-standard-64 | same count |
| Aggregate JVM-side RAM | 686 GB (VegaX, default heap) | 191 GB (capped heap) | −72% |
| One refresh at 10x models | ~67 h · ~$6,030 | ~6 h · ~$895 | −85% |
| Score parity | ground truth | identical to the last digit | exact |
| Projected annual savings | — | $400K+ | — |
- Identical scores: met. Every row matches Vega to the last digit, across every probability column, and the gate that proved it now runs as a test.
- A fraction of the cost, not just the time: met. $602.74 to $89.54 per run (−85%) and 6h 44m to 36m (−92%), on standard machines instead of highmem.
- A reusable component: met. A containerized Flex Template that any team can launch with a table, a model, and an output path, already lined up for the Tax ML launch.
- GCP-native: met. Beam, Dataflow, and BigQuery, with no new infrastructure to operate.
For the business, that turns the 10x ask from a multi-day, ~$6,000 refresh into roughly six hours and ~$895. It means facts can be refreshed far more often, and it adds up to more than $400K a year in projected savings. Some of the optimizations from this work also made it into Google Cloud's own documentation as a recommended pattern.
None of these gains came from changing the model. They came from the batching, memory settings, and data handling around it.
Acknowledgments
This project had a lot of support behind it, and I want to start with Pallav AnandStaff Machine Learning EngineerIntuit Credit KarmaThe Staff MLE I worked with on the batch scoring project.LinkedInView profile, the Staff MLE I worked with all summer. A great director makes everyone on set better, and Pallav did that for me. He introed the project well from the first week, when he walked me through PMML scoring, the existing Beam pipelines, Airflow DAGs, BigQuery sinks, so I could go find where the time was going. His staff-level insights helped me sharpen my priorities as I worked. Our Zoom sessions in the middle of parity debugging was where we decided to compare inputs feature by feature, and I took it from there. His reviews held my code to a staff-level bar, and his thinking on PMML standardization helped me see where the component fits beyond this summer. Most of all, he trusted me to own the project end to end and was always there when I wanted a second pair of eyes. With the utmost respect and sincerity, a hearteous thank you, Pallav.
Thank you to Hardik UdeshiSenior Software EngineerIntuit Credit KarmaLinkedInView profile (Senior SWE), for catching the resource exhaustion in our shared GCP resources and providing clarification on company context; to Bo Li for a thorough and patient code review; to Nicholas Pataki, Debasish Das, Dharam Pulakunta, Raj Katakam, and Maddie Daianu for the direction and support; to Donald Anthony Binns and Lejorne Leys (university recruiting team) for making the summer happen.

