Can you give a high level description of how you implement the global grid? It's really only absolute coordination of a checkerboard-like grid that I'm wrestling with. I've verified that my algorithm works, and is nicely stable. I've also scrutinized it with VisualVM, and on Linux and OS X, the threads are either in wait or running, in the pattern I'd expect. (The subgrid worker-group threads have to wait for the global grid to do its coordinating.) I'm also seeing expected use of the thread pools. However, for some reason I can't seem to scale beyond 250 users.
One complication is that my grid is for a procedurally generated world with 2^120 locations in it. This is why I generate subgrids. A degenerate case is one subgrid per user. However, these subgrids are organized in load-balanced groups, each of which has their own thread pool, caches, and locks.
Also, rollbacks are problematic, though the real problems are arguably corner cases.
Erlang might be a win because each garbage collector only has to deal with its own local memory.
EDIT: It turns out my algorithm is somewhat similar to Pikko:
http://www.erlang-factory.com/upload/presentations/297/Pikko...
One big difference, is that my algorithm doesn't move or reconfigure masts, instead it dynamically creates subgrids, which are then grouped into "workgroups" each of which is supposed to be processed by a different CPU socket. Instead of there being an API, it's more that the subgrids stop what they are doing and their information is briefly managed by the global grid code. (The procedurally generated map is rigged, so that there are many opportunities for crossing from one subgrid to another.)