This is interesting, but I think the authors don't talk enough about CPU and memory utilization. To me, the "classic" Google distributed systems architecture puts different tasks in different logical servers (doesn't matter if they're separate physical servers or not), which gives those servers more predictable memory and CPU usage, which in turn enables tighter bin-packing of jobs in the datacenter. The price they pay is needing a really, really fast in-datacenter network, but in the past they've been okay with this.
The paper proposes putting application-specific processing and memory caching on the same host, which might give the combined server less consistent CPU usage and therefore lower utilization, but will also eliminate the network hop from application server to in-memory cache. It seems intuitively reasonable to me to give up some CPU utilization in exchange for eliminating an all-to-all network connections stage, but I would like to see a real cost and speed comparison.