For example, there is a lot of caching you just have to do. Tons of it. Cache as much as you can. Memcache and Varnish are your friends, use them as much as you can. Unfortunately, if Healthcare.gov was in a very write-heavy situation, they are somewhat limited in how much caching will help.
One thing they could have done that would have saved A TON of load is not require users to sign up before giving them a list of available plans. That whole part didn't need to be database driven at all. The data could have been stored in redis and they could have used javascript to filter it based on a form. That would have been ridiculously fast. They also could have prerendered all the plan list possibilities and stored those in varnish. That also would have been ridiculously fast. My guess is millions of people just wanted to check prices and eliminating the database load for those users would have probably kept things running fast and smooth.
Slow DB queries are the enemy and you don't realize how bad they are until you are at scale. Sure, it only takes a few seconds on your local machine, but multiply that times thousands of concurrent users and your DB gets swamped. If you are using an ORM, it is MUCH harder to track down where in your code that 3 way join that scans every record is happening. Ideally you'd be using straight SQL and maybe use comments to tag a query. Also tools like newrelic might help if only because many databases don't offer great visibility of performance data.
Sharding your database is something that is probably possible, and in the case of healthcare.gov, they probably could have had totally separate infrastructure on a per-state basis that would have made scaling a lot easier than putting everything on the same database. Also, put reporting and things that aren't mission critical on a slave database. The last thing you want is a reporting job bringing down the live site in the background.
Getting good hardware with fast IO is going to save a lot of developer time required to scale things. Using fast SSD's is probably the easiest win to speed up your database. Developer time costs a lot more than hardware and giving yourself cheap headroom up front gives you breathing room on launch.
Performance testing is also something worth doing, but the tricky part is until you roll out, it is hard to know exactly where the hotspots are going to be. In this case, new user signup would be the obvious place to test, so they probably should have tested up to the limits of their servers and tried to extrapolate an expected number of users and maybe increased that by 1.5x or something to have some leeway.
On a rollout where you don't know what you are getting into user wise, being on a cloud where you can scale out fast as demand requires is something worth doing. They could have saved a lot of bad press and headache by being on the cloud initially and migrating to less hardware after the initial peak died down.
There are a ton of little things like that you have to think about if you are dealing with massive scale. I don't know if the engineers building healthcare.gov had ever dealt with something like this before, but I'm sure they're learning these lessons now.