16 ms·
Binance built a 100PB log service with Quickwit
- cletus 2y agoJust browsing the Quickwit documentation it seems like the general architecture here is to write JSON logs but stores them compressed. Is this just something like gzip compression? 20% compressed size does seem to align to ballpark estimates of JSON GZIP compression. This is what Quickwit (and this page) calls a "document": a single JSON record (just FYI). Additionally you need to store indices because this is what you actually search. Indices have a storage cost when you write them too. When I see a system like this my thoughts go to questions like: - What happens when you alter an index configuration? Or add or remove an index? - How quickly do indexes update when this happens? - What about cold storage? Data retention is another issue. Indexes have config for retention [1]. It's not immediately clear to me how document retention works, possibly from S3 expiration? So, network transfer from S3 is relatively expensive ($0.05/GB standard pricing [2] to the Internet, less to AWS regions). This will be a big factor in cost. I'm really curious to know how much all of this actually costs per PB per month. IME you almost never need to log and store this much data and there's almost no reason to ever store this much. Most logs are useless and you also have to question what the purpose is of any given log. Even if you're logging errors, you're likely to get the exact same value out of 1% sampling of logs than you are with logging everything. You might even get more value with 1% sampling because your query and monitoring might be a whole lot easier with substantially less data to deal with. Likewise, metrics tend to work just as well from sampled data. This post suggests 60 day log retention (100PB / 1.6PB daily). I would probably divide this into: 1. Metrics storage. You can get this from logs but you'll often find it useful to write it directly if you can. Getting it from logs can be error-prone (eg a log format changes, the sampling rate changes and so on); 2. Sampled data, generally for debugging. I would generally try to keep this at 10TB or less; 3. "Offline" data, which you would generally only query if you absolutely had to. This is particularly true on S3, for example, because the write costs are basically zero but the read costs are expensive. Additionally, you'd want to think about data aggregation as a lot of your logs are only useful when combined in some way [1]: https://quickwit.io/docs/overview/concepts/indexing https://quickwit.io/docs/overview/concepts/indexing [2]: https://aws.amazon.com/s3/pricing/ https://aws.amazon.com/s3/pricing/
- JackSlateur 2y agoYou have very good questions, I can only guess one answer: s3 network transfer is free for AWS services Your link[1] said: You pay for all bandwidth into and out of Amazon S3, except for the following: [...] - Data transferred from an Amazon S3 bucket to any AWS service(s) within the same AWS Region as the S3 bucket (including to a different account in the same AWS Region).
- fulmicoton 2y agoQuickwit (like Elasticsearch/Opensearch) stores you data compressed with ZSTD in a row store, builds a full text search index, and stores some of your fields in a columnar. The "compressed size" includes all of this. The high compression rate is VERY specific to logs. - What happens when you alter an index configuration? Or add or remove an index? Changing an index mapping was not available in 0.8. It is available in main and will be added in 0.9. The change only impacts new data. - Or add or remove an index? This is handled since the beginning. - What about cold storage? What makes Quickwit special is that we are reading everything is on S3. We adapted our inverted index to make it possible to read straight from S3. You might think this is crazy slow, but we typically search into TBs of data in less than a second. We have some in RAM cache too, but they are entirely optional. > 2. Sampled data, generally for debugging. I would generally try to keep this at 10TB or less; Sometimes, sampling is not possible. For instance, some of Quickwit users (including Binance) use their logs for user support too. A user might come asking details about something fishy that happened 2 months ago.
- ATsch 2y agoIt's always very amusing how all of the blockchain companies wax lyrical about all of the huge supposed benefits of blockchains and how every industry and company is missing out by not adopting them and should definitely run a hyperledger private blockchain buzzword whatever. And then, even when faced with implementing a huge, audit critical, distributed append-only store, the thing they tell us blockchains are so useful for, they just use normal database tech like the rest if us. With one centralized infrastructure where most of the transactions in the network actually take place. Who's tech stack looks suspiciously like every other financial institution. I'm so glad we're ignoring 100 years of securities law to let all of this incredible innovation happen.
- tommek4077 2y agoBinance is not a blockchain company. It is a centealized exchange. Nothing is happening on-chain unless getting coins from or to the exchange. And this has nothing tondo with them then.
- rijoja 2y ago> And then, even when faced with implementing a huge, audit critical, distributed append-only store, the thing they tell us blockchains are so useful for, they just use normal database tech like the rest if us. With one centralized infrastructure where most of the transactions in the network actually take place. Who's tech stack looks suspiciously like every other financial institution. Right but surely you must understand that the blockchain transactions are already stored in the blockchain, and what this is about is logs that might be useful for debugging purposes, and as such would be more verbose than what's required and also could contain sensitive information? Apart from that isn't it obvious that the performance requirement would make this unrealistic, with no added benefits, whatsoever. >I'm so glad we're ignoring 100 years of securities law to let all of this incredible innovation happen. Storing all these logs on a blockchain might very well (apart from being totally asinine) breach privacy regulations as well, as it might very well store sensitive data? Surely you must understand this?
- ATsch 2y ago> Surely you must understand this? Yes, I understand why blockchains are bad, have no benefits, terrible performance and are a privacy nightmare. Thanks for explaining it in more detail. And binance understands it too, that's why they're not using it (not even a private one!) despite all of their talk about how it's a revolutionary technology.
- ram_rar 2y ago>Limited Retention: Binance was retaining most logs for only a few days. Their goal was to extend this to months, requiring the storage and management of 100 PB of logs, which was prohibitively expensive and complex with their Elasticsearch setup. Just to give some perspective. The Internet Archive, as of January 2024, attests to have stored ~ 99 petabytes of data. Can someone from Binance/quickwit comment on their use case that needed log retention for months? I have rarely seen users try to access actionable _operations_ log data beyond 30 days. I wonder how much $$ can they save more by leveraging tiered storage and engs being mindful of logging.
- drak0n1c 2y agoGovernment regulators take their time and may not investigate or alert firms to identified theft, vulnerability, criminal or sanctioned country user trails for months. However, that does not protect those companies from liability. There is recent pressure and targeted prosecution from the US on Binance and CZ along this angle. They've been burned on US users getting into their international exchange, so keeping longer forensic logs helps surveil, identify, and restrict Americans better (as well as the bad guys they're not supposed to interact with).
- evdubs 2y agoLots of storage to log all of the wash trading on their platform.
- ukuina 2y agoDoes QuickWit support regex search now? The underlying store, Tantivy, already does. This is what stopped a PoC cold at an earlier project.
- PSeitz 2y agotantivy has two dictionaries FST and SSTable. We added SSTable in tantivy because it works great with object storage, while FST does not. With some metadata we can download only the required parts and not the whole dictionary. SStable does not support Regex queries, it would require a full load and scan, which would be very expensive. Your best bet currently would be to make it work with tokenizing, which is way more efficient anyways. prefix queries are supported btw
- ukuina 2y agoAre in-order queries supported? e.g., TERM1*TERM2 should return matches with those terms in that specific order.
- kristopolous 2y agoSo people don't build this out themselves? Regardless, there's some computer somewhere serving this. How do they service 1.6 PB per day? Are we talking tape backup? Disks? I've seen these mechanical arms that can pick tapes from a stack on a shelf, is that what is used? (example: https://www.osc.edu/sites/default/files/press/images/tapelibrary.jpg https://www.osc.edu/sites/default/files/press/images/tapelib...) For disks that's like ~60/day without redundancy, do they have people just constantly building out and onlining machines in some giant warehouse? I assume there's built in redundancy and someone's job to go through and replace failed units? This all sounds like it's absurdly expensive. And I'd have to assume they deal with at least 100x that scale because they have many other customers. Like what is that? 6,000 disks a day? Really? I hear these numbers of petabyte storage frequently. I think Facebook is around 5PB/daily. I've never had to deal with anything that large. Back in the colo days I saw a bunch of places but nothing like that. I'm imagining forklifts moving around pallets of shrink wrapped drives that get constantly delivered Am I missing something here? Places like AWS should run tours. It'd be like going to the mint.
- jiggawatts 2y agoThe article talks about 1.6 PB / day, which is 150 Gbps of log ingest traffic sustained. That's insane. A change of the logging protocol to a more efficient format would yield such a huge improvement that it would be much cheaper than this infrastructure engineering exercise. I suspect that everyone just assumes that the numbers represent the underling data volume, and that this cannot be decreased. Nobody seems to have heard of write amplification. Let's say you want to collect a metric. If you do this with a JSON document format you'd likely end up ingesting records that are like the following made-up example: { "timestamp": "2024-07-12T14:30:00Z", "serviceName": "user-service-ba34sd4f14", "dataCentre": "eastFoobar-1", "zone": 3, "cluster: "stamp-prd-4123", "instanceId": "instance-12345", "object": "system/foo/blargh", "metric: "errors", "units": "countPerMicroCentury", "value": 391235.23921386 } It wouldn't surprise me if in reality this was actually 10x larger. For example, just the "resource id" of something in Azure is about this size, and it's just one field of many collected by every logging system in that cloud for every record. Similarly, I've cracked open the protocols and schema formats for competing systems and found 300x or worse write amplification being the typical case. The actual data that needed to be collected was just: 391235.23921386 In a binary format that would be 4 bytes, 8 if you think that you need to draw your metric graphs with a vertical precision of a millionth of a pixel and horizontal precision of a minute because you can't afford the exabytes of storage a higher collection frequency would require. If you collect 4 bytes per metric in an array and record the start timestamp and the interval, you don't even need a timestamp per entry, just one per thousand or whatever. For a metric collected every second that's just 10 MB per month before compression. Most metrics change slowly or not at all and would compress down to mere kilobytes.
- AJSDfljff 2y agoUnfortunate the interesting part is missing. Its not hard at all to scale to PB. Junk your data based on time, scale horizontally. When you can scale horizontally it doesn't matter how much it is. Elastic is not something i would use for scaling horizontally basic logs, i would use it for live data which i need live with little latency or if i do constantly a lot of log analysis live again. Did Binance really needed elastic or did they just start pushing everything into elastic without every looking left and right? Did they do any log processing and cleanup before?
- fulmicoton 2y agoThis is their application logs. They need to search into it in a comfortable manner. They went for a search engine with Elasticsearch at first, and Quickwit after that because even after restriction the search on a tag and a time window "grepping" was not a viable option.
- AJSDfljff 2y agoWould be curious what they are searching exactly. At this size and cost, aligning what you log should save a lot of money.
- BiteCode_dev 2y agoFinancial institutions have to log a lot just to comply with regulations, including every user activity and every money flow. On an exchange that does billions of operation per seconds, often with bots, that's a lot.
- AJSDfljff 2y agoYes but audit requirements doesn't mean you need to be able to search everything very fast. Binance might not have a 24/7 constant load, there might be plenty of time to compact and write audit data away at lower load while leveraging existing infrastructure. Or extracting audit logging into binary format like protobuff and writing it away highly optimized.
- randomtoast 2y agoI wonder how much their setup costs. Naively, if one were to simply feed 100 PB into Google BigQuery without any further engineering efforts, it would cost about 3 million USD per month.
- AJSDfljff 2y agoGood question. I thought it would be a no brainer to put it on s3 or similiar but thats already way to expensive at 2m/month without api requests. Backplace storage pods are an initial investment of 5 Million, thats probably the best bet you could do and on that savings level, having 1-3 good people dedicated to this is probably still cheaper. But you could / should start talking to the big cloud providers to see if they are flexible enough going lower on the price. I have seen enough companies, including big ones, being absolut shitty in optimizing these types of things. At this level of data, i would optimize everyting including encoding, date format etc. But i said it in my other comment: the interesting questions are not answered :D
- orf 2y agoThe compressed size is 20pb, so it’s about 500k per month in S3 fees
- francoismassot 2y agoIndeed. They benefit from a discount, but we don't know the discount figure. To further reduce the storage costs, you can use S3 Storage Classes or cheaper object storage like Alibaba for longer retention. Quickwit does not handle that, so you need to handle this yourself, though.
- bearjaws 2y ago[flagged]
- martijnvds 2y agoExcept for AWS ;)
- deleted 2y ago[deleted]
- endorphine 2y agoWhat would you use for storing and querying long-term audit logs (e.g. 6 months retention), which should be searchable with subsecond latency and would serve 10k writes per second? AFAICT this system feels like a decent choice. Alternatives?
- jakjak123 2y ago10k audit logs per sec? I think we have different definitions of audit logs.
- jjordan 2y agoNATS?
- packetlost 2y agoNATS doesn't really have advanced query features though. It has a lot of really nice things, but advanced querying isn't one of them. Not to mention I don't know if NATS does well with large datasets, does it have sharding capability for it's KV and object stores?
- Zambyte 2y agoI use NATS at work, and I have had the privilege to speak with some of the folks at Synadia about this stuff. Re: advanced querying: the recommended way to do this is to build an index out of band (like Redis (or a fork) or SQLite or something) that references the stored messages by sequence number. By doing that, your index is just this ephemeral thing that can be dynamically built to exactly optimize for the queries you're using it for. Re: sharding: no, it doesn't support simple sharding. You can achieve sharding by standing up multiple NATS instances, and making a new stream (KV and object store are also just streams) on each instance, and capture some subset of the stream on each instance. The client (or perhaps a service querying on behalf of the client) would have to me smart enough to be able to mux the sources together.
- packetlost 2y ago
- kstrauser 2y agoHow? If Binance had a trillion transactions, that’s 100KB per transaction. What all are they logging?
- tommek4077 2y agoHigh frequency traders are making hundreds of billions of orders per day. And there are many bigger and smaller players.
- francoismassot 2y agoThey have 181 trillion logs
- kstrauser 2y agoBut of what? What has Binance done 181 trillion times? Obviously they have. I don’t think they’re throwing away money for logs they don’t generate or need. I just can’t imagine the scope of it. That is, I know this is a failing of my imagination, not their engineering decisions. I’d love to fill in my knowledge gaps.
- xboxnolifes 2y agoThey are application logs, so probably nearly every click on their website.
- robxorb 2y agoIf it's 181 trillion each year, it's only 6 million per second. There's a thousand milliseconds in each second so Binance would need only several thousand high frequency traders creating, and adjusting orders, through their API, to end up with those logs. Binance has hundreds of trading pairs available so a handful on each pair average would add up.
- askl 2y agoDon't forget the logs produced by the logging infrastructure.
- KaiserPro 2y agoA word of caution here: This is very impressive, but almost entirely wrong for your organisation. Most log messages are useless 99.99% of the time. Best likely outcome is that its turned into a metric. The once in the blue moon outcome is that it tells you what went wrong when something crashed. Before you get to shipping _petabytes_ of logs, you really need to start thinking in metrics. Yes, you should log errors, you should also make sure they are stored centrally and are searchable. But logs shouldn't be your primary source of data, metrics should be. things like connection time, upstream service count, memory usage, transactions a second, failed transactions, upsteam/downstream end point health should all be metrics emitted by your app(or hosting layer), directly. Don't try and derive it from structured logs. Its fragile, slow and fucking expensive. comparing, cutting and slicing metrics across processes or even services is simple, with logs its not.
- ryukoposting 2y agoMetrics are useful when you know what to measure, which implies that you already have a good idea for what can go wrong. If your entire product exists in some cloud servers that you fully control, that's probably feasible. Binance probably could have done something more elegant than storing extraordinary amounts of logs. However, if you're selling a physical product, and/or a service that integrates deeply with third party products/services, it becomes a lot more difficult to determine what's even worth measuring. A conservative approach to metrics collection will limit the usefulness of the metrics, for obvious reasons. A "kitchen sink" approach will take you right back to the same "data volume" problem you had with logs, but now your developers have to deal with more friction when creating diagnostics. Neither extreme is desirable, and finding the middle ground would require information that you simply don't have. On a related note, one approach I've found useful (at a certain scale) is to shove metrics inside of the logs themselves. Put a machine-readable suffix on your human-readable log messages. The resulting system requires no more infrastructure than what your logs are already using, and you get a reliable timeline of when certain metrics appear vs. when certain log messages appear.
- temporarely 2y agoAny system has a 'natural set' of metrics. And metrics are not about "what [went] wrong" rather system health. So Metrics -> Alert -> Log Diagnostics.
- sebstefan 2y ago> On a given high-throughput Kafka topic, this figure goes up to 11 MB/s per vCPU. There's got to be 2x to 10x improvement to be made there, no? No way CPU is the limitation these days and even bad hard drives will support 50+mB/s write speeds.
- ddorian43 2y agoBuilding inverted index is very CPU intensive.
- fulmicoton 2y agoBuilding an inverted index is actually very cpu intensive. I think we are the fastest on that (if someone knows something faster than tantivy at indexing I am interested). I'd be really surprised if you can make a 10x improvement here.
- ZeroCool2u 2y agoThere was a time at the beginning of the pandemic where my team was asked to build a full text search engine on top of a bunch of SharePoint sites in under 2 weeks and with frustratingly severe infrastructure constraints, (No cloud services, single box on prem for processing, among other things), and we did and it served its purpose for a few years. Absolutely no one should emulate what we built, but it was an interesting puzzle to work on and we were able to cut through a lot of bureaucracy quickly that had held us back for a few years wrt accessing the sensitive data they needed to search. But I was always looking for other options for rebuilding the service within those constraints and found Quickwit when it was under active development. I really admire their work ethic and their engineering. Beautifully simple software that tends to Just Work™. It's also one of the first projects that made me really understand people's appreciation for Rust as well outside of just loving Cargo.
- totaa 2y agoI don't know what brings me more happiness in this career. Building systems with no political constraints, or building something that's functional with severe restraints.
- fulmicoton 2y agoThank you for the kind word @ZeroCool2u ! :)
- hanniabu 2y ago> we were able to cut through a lot of bureaucracy quickly that had held us back for a few years wrt accessing the sensitive data they needed to search Doesn't sound like a benefit for your users
- shortrounddev2 2y agoIn what way?
- hanniabu 2y agoBypassing protections for accessing sensitive data...
- elchief 2y agomaybe drop the log level from debug to info...
- piterrro 2y agoReminds me of the time Coinbase paid DataDog $65M for storing logs[1] [1] https://thenewstack.io/datadogs-65m-bill-and-why-developers-should-care/ https://thenewstack.io/datadogs-65m-bill-and-why-developers-...
- inssein 2y agoWhy do so many companies insist on shipping their logs via Kafka? I can't imagine deliverability semantics are necessary with logs, and if they are, they shouldn't be in your logs?
- jcgrillo 2y agoKafka is a big dumb pipe that moves the bytes real fast, it's ideal for shipping logs. It accepts huge volumes of tiny writes without breaking a sweat, which is exactly what you want--get the logs off the box ASAP and persisted somewhere else durably (e.g. replicated).
- ecnahc515 2y agoIn log shipping cases it’s good as a buffer so you can batch writes to the underlying SIEM. This prevents tons of small API calls with a few hundred or thousand log lines each. Instead Kafka will take all the small calls and the SIEM can subscribe and turn them into much larger batches to write to the underlying storage (eg S3).
- bushbaba 2y agoDon’t forget about all the added cost. never got it as many shops can tolerate data loss for their melt data. So long as it’s collected 99.9% of the time it’s good enough.
- mdaniel 2y agoMy experience has been a mixture of "when all you have is a hammer ..." and Pointy Haired Bosses LOVE kafka, and tend to default to it because it's what all their Pointy Haired Boss friends are using In a more generous take, using some buffered ingest does help with not having to choose between a c500.128xl ingest machine and dropping messages, but I would never advocate for standing up kafka just for log buffering
- inssein 2y agoat that point you are likely slowing down your applications - I think a basic OpenTelemetry collector mostly solves this, and if you go beyond the available buffer there, then dropping it is the appropriate choice for application logs.
- RIMR 2y agoI am having trouble understand how any organization could ever need a collection of logs larger than the size of the entire Internet Archive. 100PB is staggering, and the idea of filling that with logs, while entirely possible, just seems completely useless given the cost of managing that kind of data. This is on a technical level quite impressive though, don't get me wrong, I just don't understand the use case.
- tommek4077 2y agoThese are order and trade logs probably. You want to have them and you need them for auditing. Binance wants to be more professional in that way probably. HFT is making billions of orders per day per trader.
- jcgrillo 2y agoOK, so let's do some napkin math... I'm guessing something like this is the information you might want to log: user ID: 128bits timestamp: 96bits ip address: 32bits coin type: idk 32bits? how many fake internet money types can there be? price: 32bits quantity: 32bits So total we have 352bits. Now let's double it for teh lulz so 704bits wtf not. You know what fuck it let's just round up to 1024bits. Each trade is 128bytes why not, that's a nice number. That means 200Pb--2e17 bytes mind you--is enough to store 1.5625e16 trades. If all the traders are doing 1e9 trades/day, and we assume this dataset is 13mo of data, that means there are 38772 HFT traders all simultaneously making 11574 trades per second.. That seems like a lot.. In other words, that means Binance is processing 448.75 million orders per second.. Are they though? EDIT: No, indeed some googling indicates they claim they can process something like 1.4 million TPS. But I'd hazard a guess the actual figure on average is less.. EDIT: err sorry, shoulda been 100Pb. Divide all those numbers by two. Still two orders of magnitude worth of absurd.
- RIMR 2y agoThe only thing I can think of is that they are collecting every single line of log data from every single production server with absolutely zero expiration so that they can backtrack any future attack with precision, maybe even finding the original breach. That's the only actual use case I can think of for something like this, which makes sense for a cryptocurrency exchange that is certainly expecting to get hacked at some point.