• Hacker News
  • new|
  • comments|
  • show|
  • ask|
  • jobs|
  • jandrewrogers 7 minutes

    > Unlike C++, which allows implicit memory allocations and complex copy constructors

    It is trivial to ensure this doesn't happen at compile-time. If you are doing DMA as in the article these guarantees are required.

  • kilroy123 4 hours

    Their simulation is fantastic:

    https://sim.tigerbeetle.com

  • g_delgado14 1 hours

    Is it just me or does it not seem like this whole article was llm generated?

    The two "sources" are fake links that lead to 404s

    lmz 10 minutes

    Look at the site. A post every 2 days? Must be a very productive person. Or just generated.

  • SPascareli13 4 hours

    Amazing writing, very accessible too.

    So batching requests is always something I think should increase performance by a lot, but most server implementations make this pretty difficult, but the thing I struggle the most to understand is how to keep the latency down if you have multiple clients request all batched together? The total amount of latency for all clients is always the latency for the slowest.

    whizzter 4 hours

    I think you design in layers, frontends that work as clients to TigerBeetle for work in batches (as mentioned by the sibling comment), but the whole idea of removing latency differences by removing unpredictability means that you don't get the jitter of latency differences that can cause backing up in normal scenarios.

    jorangreef 4 hours

    If you give the TigerBeetle client a single transfer, it sends it off immediately to the cluster. There's no delay. No Nagle!

    But if your application then creates another transfer against the client, and another, while the first request is inflight, then the client will autobatch under the hood and send these off as a batch when the first request returns.

    You get this sweetspot then between latency and throughput. And your latency is not spiking as your load increases, since your throughput is now able to keep up.

  • teabee89 3 hours

    "By [...] utilizing a single-threaded execution loop, TigerBeetle aligns its software architecture perfectly with the physical realities of modern hardware."

    Can someone explain why single-threaded execution loop is more aligned with the physical realities of modern hardware ?

    to_ziegler 3 hours

    Tobi here from TB. Great question, and it's important to be nuanced here.

    It really depends on the problem space. For example, many OLAP workloads (analytical) contain large amounts of parallelizable work, then multi-core execution is absolutely the way to go. That also aligns well with the direction CPU technology is taking, with core counts continuing to increase.

    For the transactional workloads we see at TigerBeetle, and in other transactional systems I've worked with, the picture is quite different. We see a lot of read-modify-write operations combined with a power-law distribution of the data.

    Take a simple banking example: some accounts, such as those belonging to large online retailers, see much more activity than the average individual account. You might have 80 - 90% of transfers touching a relatively small number of these hot accounts.

    Operations on the same account must be serialized to preserve correctness. That means this part of the workload cannot be meaningfully parallelized. In fact, attempting to parallelize it can make performance worse because of lock contention and coordination overhead, something the "Universal Scalability Law" captures quite well (but is also easy to test out yourself with a simple experiment).

    Instead, we focus on batched execution. We carefully structure execution to make effective use of CPU caches and efficient algorithms, so that a single batch can be processed extremely efficiently without any coordination. Batch execution also allows to amortize I/O and replication.

    That being said, there are areas where we could use multi-threading (e.g. compaction) that are not on the hot execution path.

    hansvm 2 hours

    Ha, I was explaining this just yesterday to a few people at $WORK (I'm not at TB, though we do use a bunch of Zig in prod).

    "Multiprocessing" (multiple CPU cores independently executing and only able to coordinate via some message-passing system -- a definition which encompasses both multi-core CPUs and horizontally scaled distributed systems, contrasting slightly with the normal definition) is challenging for a few reasons.

    Firstly, the details are an open math question, but I'll blindly state that some problems aren't amenable to parallelization. I.e., no algorithm can meaningfully improve performance via parallelization no matter the implementation. Think through how you would more quickly compute hash(hash(hash(hash(...)))) for example. The serial dependency makes things challenging. That isn't too dissimilar to the problem TB faces.

    Secondly, message passing is expensive. If the only way two CPU cores can coordinate is through a multi-level cache, at best you're incurring ~tens of nanoseconds of latency per message. Contrast that with a base rate of 512 bytes processed per nanosecond with enough attention to detail on typical modern server hardware (4 pipelined AVX512 instructions at 2GHz). Messages are several orders of magnitude worse than your normal work, so if you need very many of them then you're hosed from a performance perspective (worse with longer delays, like networked computers). Even very parallelizable problems at an abstract level can suffer performance losses by trying to add even one extra core. This blog post [0] doesn't perfectly capture the idea, but it's close (and a fun read regardless).

    Thirdly, message passing is an insanely complicated programming abstraction to reason about. My first two points were more about what TB was saying -- realities of modern hardware -- but the programmer experience is important too (even if you don't believe that post-2020, the LLM experience doesn't differ much from the human experience; bad code begets more bad code, slowly). The core mechanism for correctness in most software is being able to reason about "this thing is true, therefore that thing is true" and iterating. You rely on invariants like "this is sorted" to build other working theories. The invariants in multiprocessing code are much more nebulous and less amenable to accidental discovery, also less amenable to being able to build or compose them into other stronger invariants as you add code. The main reason for that is that you know almost nothing about the relative order of those messages with respect to each processer's view of which instructions happened when (and for purely multi-core "multiprocessing" the story is even worse; while my description of message-passing being the core primitive is correct, that's not what's exposed to you as a programmer, and different memory models can have even weirder interleavings than your code would naively suggest -- i.e., your code is being decomposed into smaller subunits than even a single assembly instruction, and the message passing happens at that level). The combinatorial explosion (an exponential explosion really, but big numbers either way) of states you might be interacting with makes it very difficult to understand _anything_ about the system you're examining. That's why you see a handful of primitives used over and over -- if you can decompose your problem into a parallel map plus an associative reduce then you can probably figure out some way to make it better through parallelization (not always, especially if the framework is too generic, see the linked blog post [0] if you weren't enticed to read it previously). If you can't decompose it into know primitives then it's an open research problem every time.

    The crux of that third point (and we could definitely add more explanation and additional problems) is that there's a huge cost to multiprocessing. You have to be buying something substantial to even want to reach for it, else you have to be in one of the "easy" problem spaces where somebody else has done the hard work (e.g., stateless webservers).

    TB isn't that. Their whole raison d'être is state management, and not in a way that's easily amenable to parallelization.

    [0] https://adamdrake.com/command-line-tools-can-be-235x-faster-...

  • hoppp 4 hours

    I really wish they turned it into a dependency or database framework where users could define their own business logic to swap out the double entry accounting, while reusing all the system architecture and networking features, consensus etc.

    Sort of like a new paradigm where opinionated custom databases could be created with arbitrary entry logic built on this stack.

    mrkeen 4 hours

    Emphasis on: opinionated custom databases. One database might be SQL-based for periodic report generation. Another database might be a noSQL key-val store designed to effortless grow with the number of end-users. General-purpose programming languages already cater to the 'own business logic' part - it's their whole job. What's left? System architecture, networking features, consensus. That's Kafka.

    high_na_euv 4 hours

    llvm for databases?

    mgrandl 3 hours

    Isn’t that literally what Turso said they wanted to be?

    Ygg2 4 hours

    I'm pretty sure they manage to get that level of performance and reliability, since they have a very limited schema.

    I'm pretty sure you can't just do a precise 128 byte align, if one of the element is an image blob or varchar(1000).

    jorangreef 4 hours

    No, the LSM in TB is generalizable at comptime to any combination of power of two sized key/value tuples.

    jorangreef 4 hours

    This was always the plan, and if you look closer at VSR and the state machine interface you’ll see it’s already pluggable. We just haven’t packaged it. (We’re dogfooding our first few internal “CustomBeetles” before we package and document.)

    ksec 3 hours

    In about 5 years time the trend will be "Just use Beetles" for everything :P

    Seriously Can't wait to see this.

    Congrats on launching [1] Tigerbeetle Cloud. Should have submitted that as well but I thought this was more interesting. May be another time.

    [1] https://tigerbeetle.com/cloud

    jorangreef 2 hours

    Ah thank you! (and for posting!)

    That’s the dream. To serve the world’s transactions and data (with a whole lotta beetles!). We’re working to make it reality.

    Appreciate the congrats! You make today a double whammy! :P

    hoppp 4 hours

    That's cool! I will be then probably taking a closer look!

    jorangreef 4 hours

    Thanks! Watch IronBeetle too on Twitch if you’d like to go really deep.

  • jorangreef 5 hours

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

    Alien1Being 3 hours

    Great and interesting work !

    dustbunny 3 hours

    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?"

    jorangreef 3 hours

    Let’s answer that when we’re there!

    ofiryanai 4 hours

    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?

    jorangreef 4 hours

    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.)

    jandrewrogers 14 minutes

    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.

    joosd 4 hours

    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?

    jorangreef 3 hours

    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!

    emj 4 hours

    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.

    to_ziegler 3 hours

    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.

    hawk_ 3 hours

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

    to_ziegler 3 hours

    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.

    rfw300 3 hours

    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.

    jorangreef 3 hours

    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.

    sashank_1509 2 hours

    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.

    27183 2 hours

    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.

    jorangreef 2 hours

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