Joran from TigerBeetle here! I created TB. Happy to answer questions!

Hey, very inspiring article, Redis engineer here. How do you work with static allocation on variable query structure, and row counts that can explode depending on the data shape?

And isn't there a benefit for small allocations on advanced memory allocations that you can't leverage if all is working in big page allocations? Do you implement memory allocations from scratch or leveraging existing allocator on top of these memory blocks strategy somehow?

Thanks! We use streaming data structures.

For example, if you take a look at our LSM compaction, regardless of the table size, we compact at the 512 KiB block granularity, and everything is streaming.

The same principle applies everywhere.

In our experience writing TigerStyle (and for all our internal code and tooling, not only TB as DBMS), we’ve never had a scenario where static allocation was not applicable or didn’t produce a better design.

You also tend to become more memory efficient, not less. Again, since you’re streaming. (You’re not allocating a massive buffer, just because a file is multi-GiB.)

I've long used similar static allocation models in analytical database kernels. These are even more susceptible to widely varying demands on memory. There are few practical limitations or caveats to this allocation model and it has strong advantages for both robustness and performance engineering. Memory organized as pages is compatible with small allocations, and is more or less how classic allocators work.

The runtime allocation is type-aware, workload-aware, and schedule-aware. The last is most important. There are two places that can act as a sink for heavy memory demands: storage (i.e. paging to disk) and network (e.g. streaming results). These have their own limitations because I/O bandwidth is finite. Effectively, your allocation rate is equivalent to available I/O bandwidth.

The most powerful lever you have to manage this is total control of the schedule. Demands on memory are created by a set of operations or queries visible to the software. The scheduler doesn't incrementally execute these operations randomly, it continuously selects execution based on the availability of memory or bandwidth to absorb the allocation demand of the operation. The scheduler has the ability to control the allocation rate to instantaneously match availability.

This is essentially the very old idea of "optical buffering" -- treating fiber optic cables as RAM -- taken to its logical architectural conclusion.

The caveat is that this requires direct I/O in userspace, which places limits on software architecture. But if you care about performance, you'd be using this type of software architecture regardless.

If one has a TigerBeetle cluster with high inter node latency, are there any easy wins to lower the latency of the whole cluster left? My head hurts when I think of latency in large clusters, so your work this year with latency was inspiring.

EDIT: I guess part of the question is about the problems with clusters with >130ms latency and if there are challenges you consider easy.

Hi, Tobi here from TB. Great question! Generally, 130 ms of network latency is challenging, and there often isn't an easy way around it as you're ultimately constrained by the speed of light (e.g. cross region deployments).

That said, network latency usually follows a distribution. For example, the median might be 130 ms while p99 is 200 ms. So one important goal is to avoid being affected by the high-latency tail.

In consensus and replication systems such as TigerBeetle, you can reduce the impact quite a bit by taking advantage of the fact that you only need a quorum. We have six replicas, and under normal operation we only need acknowledgements from three (including the primary, since we use flexible quorums). That means the primary only has to wait for the two fastest replicas to respond. This is very effective at reducing tail latency.

Then, to get as close as possible to speed-of-light latency, you want to avoid adding unnecessary latency inside the system itself. We've done quite a few algorithmic optimizations there over the past year. For example, introducing radix sort and tournament trees to make CPU processing more efficient.

Is this flexible quorum something that deviates from VSR? I thought VSR requires majority of the nodes to form quorum.

It's an insight from Heidi Howard et al. that came out after VSR: https://arxiv.org/abs/1608.06696 and can be applied to VSR (and others).

The basic idea is pretty simple. In VSR, there are two main phases:

1. Leader election

2. Normal replication / request processing

Before Heidi Howard’s insight, these two phases typically used the same quorum size - for example, 4 out of 6 replicas.

The key observation was that the two phases can actually use different quorum sizes, as long as the relevant quorums still intersect.

With 6 replicas, we could use a quorum of 4 for view change and a quorum of 3 for normal processing, because 4+3>6. This guarantees that every view-change quorum intersects every processing quorum. Therefore, if an operation was committed by a processing quorum, at least one replica participating in the subsequent view change knows about that operation. Combined with the protocol's view-change/log-selection rules, this ensures that committed operations are preserved when the new leader takes over.

If this interests you, Heidi gave a talk about this at systems distributed: https://youtu.be/P0cAG-RM1_c which will be released soon.

Hey Joran, awesome stuff. I see a TB post every now and then and it seems like such an interesting problem space to work in. I'm all the way at the other end of the stack, most software we write day-to-day is in JVM-based languages where you don't have to think about these things at all (we get by with ms latencies instead of ns). Reading this post inspires me explore low-level engineering more and I figure zig or rust would be a good place to start.

What are some interesting problems or things you can think of to work on that would give someone new a nice amount of exposure to this kind of programming?

Hey Joos, ah appreciated!

I'd check out tigerstyle.dev, pick up Zig, and then make an HTTP server or file format parser. Those are great ways to learn and experience this kind of programming. At some point, you start to realize that it's just easier to build API services in this way.

But for sure, you can learn so much in JVM-based languages. They make you appreciate low-level techniques all the more!

How has AI changed the way that TigerBeetle does software engineering? Given the project’s idiosyncratic language/memory allocation choices, it’s an interesting data point how well the frontier models work for you guys.

They really don’t work for us. The quality is just so poor.

We still write, read (and have an independent engineer review) each line of code by hand.

We go faster like that, but, most of all, it’s the guarantee we make to our users, also to continue to invest in our own understanding, because second order that’s valuable for the kind of high performance safety work we do.

Long term, I’m sure LLMs will improve, but right now they’re just not there.

That sounds great, I wish I had a job like that. I just wrangle agents for everything now. I think an added benefit of what you’re doing is it is more fun and this the humans who are actually building the product are more motivated to continue giving their best. With Agents it’s common to just say good enough and move on.

Thank you for this candid answer. In the current climate of people breathlessly, hyperbolically jabbering about how AI is "revolutionizing everything" it's extremely refreshing to hear this honest, measured statement.

Ah it’s a pleasure. It’s our experience, and happy to share.

Great and interesting work !

What use cases is Tiger Beetles architecture not a good fit for?

When tiger beetle becomes pluggable to different use cases, how should I think about "do I want tiger beetle?"

Let’s answer that when we’re there!