The problem with conventional databases
paulbuchheit.blogspot.com
paulbuchheit.blogspot.com
I don't imagine that deferred logging is a big deal, though, because log writes are by definition sequential. As Paul pointed out, sequential writes just aren't that slow. You can build an array of 62 commodity drives for maybe $4000.
The huge win is if you can append many transactions to the log in each rotation. To do that you have to gather up many updates per disk operation. So deferred logging is critical.
I suspect the reason PostgreSQL doesn't really support delayed log flush is that they are thinking about ACID transactions, where you really need the data to be on disk immediately. A more technical issue is that the log data must be on disk before the corresponding permanent data (otherwise crash recovery will break), and I suspect postgresql.conf's "fsync" option has the effect of not fsync()ing the log at all, which indeed would cause permanent corruption after a crash.
Indeed, fsync = off just means that the WAL isn't fsync'd at all, which can cause permanent corruption after a crash.
PostgreSQL does support a "deferred logging" mode, in which one or more transactions can avoid fsync'ing the WAL without risking data corruption -- the only risk is that those particular transactions might not be durable if the system crashes before the next fsync. This allows you to mix must-be-durable transactions with more transient ones, which is a nice feature.
You might imagine that the disk would write-cache only an amount of data that it could write to the surface with the energy stored in its capacitors after it detected a power failure. But this is not the way disks work. Typical disk specs explicitly say that the contents of the write-cache may be lost if the power fails.
You may be thinking of "tagged queuing", in which the o/s can issue concurrent operations to the disk, and the disk chooses the order in which to apply them to the surface, and tells the o/s as each completes so the DB knows which transaction can now continue. That's a good idea if there are concurrent transactions and the DB is basically doing writes to random disk positions. In the log-append case we're talking about, tagged queuing is only going to make a difference if we hand lots of appends to the disk at the same time. In that specialized situation it's somewhat faster to issue a single big disk write. You need to defer log flushes in either case to get good performance.
That's exactly what I assumed, at least for high-end disks. Any idea why they don't do that? It seems like a pretty trivial hardware feature that would save an awful lot of software complexity.
Such a feature would anyway be hard or impossible to use as part of a design to get fast writes and crash recovery. Crash recovery usually depends on constraints on the order writes were applied to the disk surface -- for example that all the log blocks were on the surface before any of the B-Tree blocks. Or (for FFS) that an i-node initialization goes to the surface before the new directory entry during a creat(). Drives that just provide write caching don't guarantee any ordering (much of the point of write-caching is to change the order of writes), and don't tell the o/s which writes have actually completed. So the write-order invariants that crash recovery depends on won't hold with write-caching. That's why tagged command queuing is popular in high-end systems: TCQ lets the drive re-order concurrent writes, but tells the o/s when each completes, so for example a DB can wait for the log writes to reach the surface before starting the B-Tree writes.
In our case, perhaps a pure log-structured DB could use a disk write-cache. Crash recovery could scan the whole disk (or some guess about the tail of the log) looking for records that were written, and use the largest complete prefix of the log. But we would not be able to use the disk for anything with a more traditional crash recovery design -- for example we probably could not store our log in a file system! Perhaps we could tell the disk to write-cache our data, but not the file system's meta-data. On the other hand perhaps we'd want to write the log to the raw disk anyway, since we don't want to be slowed down by the file system adding block numbers to the i-node whenever we append the log.
http://www.postgresql.org/docs/8.2/interactive/runtime-config-wal.html
Databases SHOULD do these things, but I haven't seen any evidence that the popular ones do. Feel free to post actual measurements for updates per second from your db.
I'm mostly an academic, so I can tell you a lot more about how things are supposed to work than about how they actually do. The only high-traffic DBMS with real users to which I can easily get access is MS-SQL, which doesn't inspire me to any leaps of faith. But once I finish building my new desktop system (with spiffy RAID array) I'll rig up some benchmarks on MySQL and PostgreSQL.
Edit: and yes, the irony is noted that I'm telling a Googler that he's not working with enough data.
I really think this concept works though for 8/10 applications which are using databases because, well, you've always used databases. Memory, save changes to flat files, read it all in during startup, it's really an elegant way to go about it.
Same setup but --innodb_flush_log_at_trx_commit=0, a million transactions in 200 seconds, or 5000/second.
I don't know if InnoDB wrote its B-Tree to disk during the 200 seconds. Same performance even with InnoDB's buffer pool size set to 200 KB with --innodb_buffer_pool_size=200000.
The MySQL documentation claims that this configuration does crash-recovery correctly, though you may lose the last second's worth of transactions.
The problem is that the Ruby / Python / PHP stacks typically have one process per request, and they do not share memory. When you use memcached in such environment, it is a separate process and communication with it is rather slow - it involves marshaling the data, which would not be necessary if all requests were just threads in a single address space.You don't have this problem if you use Java/C++.
I have used with some success a huge mmaped file as persistent memory to store my data. It's not exactly using disc as sequential device, that depends on the application and data locality. There are just a few tools available to support that scheme.
How do you handle features that rely on frequent UPDATE statements, like say a hitcounter? Is UPDATE LOW PRIORITY (MySQL) enough, or should you work out some application-specific caching/batching mechanism and perform all your updates at once?
Shouldn't that be "it's faster"?
Excellent post otherwise, btw.
If there's any silver bullet to scaling most web apps it's properly utilizing gobs of memory. Pushing bytes from memory to your NIC is a pretty damn efficient operation.
It does? Where can I get that?