This was based on yearly RAM cost. I worked on a log analysis system. In the beginning, it ran on one computer. That one program would read logs, generate fleet-wide analytics, and serve those out of RAM. Eventually we wanted to run on more than one computer, for both scaling and reliability reasons. The design of the system allowed us to basically do the same serving with many replicas, so for a long time we just ran a fleet of replicas that still had all the capabilities of the monolith. The change we made was to move the aggregation stage to a new dedicated program, that we only ran 3 copies of worldwide. (A man with 2 replicas never knows which one is broken, they say.) This meant that the mappers became significantly lighter RAM-wise, they just existed to use as much CPU time as they could, and then we had 3 beefy reducer replicas to aggregate and serve data.
This was a very easy change; we just made a new main.go and an RPC to send the data to aggregate. The system was designed internally to be logically isolated across that boundary, so we just stuck in the RPC and then the other side of the boundary could be another data center.
In the end, I think we saved a few terabytes of RAM-years. Not a big deal, but it was something.