LinkedIn adopts protocol buffers and reduces latency up to 60%
infoq.com
infoq.com
Feels that their problem may just be they went a bit too far on the microservices band wagon.
Big Tech companies will often conflate micro service with large distributed system.
These services are by no means at all micro.
Then again, GitHub still uses AWS IIRC.
I'm sorry Dave, I'm afraid I can't do that.
What makes a service micro vs macro?
cd hello_world
npm install
* watches text scroll for five hours
The real semantic fun begins when you have monolithic codebase/deployments that could do everything, but each instance gets assigned a specialized role it serves to its peers _ (I think I've seen someonevdescribing that approach somewhere). I'd consider that micro if there are more than a handful of "roles", perhaps with some qualifier, but it would hopefully be easy to agree than no description is entirely wrong.
(and the actual me who has never reached beyond consumer side in conferences sure hopes that conferences would resist unless a less ridiculous name was found)
(so, yes, I'll take your nope for my whynotting)
Now, if you are launching a search engine, you need to name it wisely for when people verb it. As a brand you probably want this to happen. Duck Duck Go was not very wise.
So number of endpoints can quickly grow to a big number.
So if you have an API with 10 endpoints, and one of them changes 10 times… you now have 100 endpoints.
It really reaches absurdity when the first endpoint is on its tenth iteration (the others haven’t changed) and now you’re serving ten duplicate endpoints per version, or 100 total endpoints where 90 of them are duplicate of themselves.
(and really the problem isn't basic CRUD endpoints, it's the ones with complex logic and structure where what's being built isn't necessarily the same thing over time.)
It's one thing when v2 and v5 are the latest, but if someone else comes through later and wants to bolt on a feature to a service that is trying to talk to v2/v5 when v3/v9 are the latest, you have to go back and look up a map of which endpoints are contemporary if you want to add a third call (v2/v5/v2) that is supposed to work together.
This can be done via swagger/etc but you are essentially just rebuilding that service versioning with an opaque "api publish date" built over the top.
But maybe the only changes between API v5 and v10 were to 5% of the endpoints. But the other 95% of the endpoints got a new version number too. That way people can refer to “API v10” instead of “Here’s a table with a different version number for all 19,000 endpoints we’re consuming in this update on our micro service”.
It’s an organizational communication thing, not a technical thing. The “API v10” implies a singular contract. Otherwise how do you communicate different version numbers for 19,000 endpoints without major miscommunications? You couldn’t even reasonably double check the spreadsheet sent between teams. Instead it’s “just make sure to use v10”. Communication is clear this way.
Obviously this method has pros and cons, I’ve explained the pros. Also this is why chaos engineering can help by intermittently black-holing old API endpoints to encourage teams to move to new ones and finally remove the old versions entirely so you don’t ever get to 19,000 endpoints, which is the real problem.
API endpoints is is almost as weird of a metric as LoC. It does tell you something, but in a way that can be misleading.
It would be a nightmare to consume something like /api/myservice/endpoint/v2. Needing v2 of the create endpoint but only v5 of the update? That would be ugly to try and work against. And actually there is no guarantee versions are even behavior compatible (although it would be stupid for it to wander too far). There can be cases where response objects don’t convey some info you need in some versions etc.
I was thinking of service as being the unit of "API" here rather than an API consisting of multiple services, "each service provides its own API" is how I was thinking of it. But I can see the usage of saying "this is our [overall] public/internal APIs" too. And I agree /api/v2/myservice would be a bit much if every service moved the global version counter every time a single endpoint was changed lol
(although I suppose you could make an argument for "timestamp" as a "point in time" API version, if you version the API globally. Sounds like it would cause friction as services try to roll out updates, but it's notionally possible at least.)
Because of that I prefer to version each service instead of versioning whole API, but both of those strategies have pros and cons.
If you're on v37 of a service and your forced to continue to support v1 (and 35 others) there's a problem somewhere.
If it's internal APIs, they need to get on top of deprecating and removing older ones. This is one of the key points of Google's SWE book (at least the first part) and the benefits of a monorepo; if you change an API in a backwards incompatible way, you're also responsible for making sure every consumer is updated accordingly. If you don't, either you're left maintaining the now deprecated API, or you're forcing however much teams to stop what they're doing and put time into a change that you decided needed to happen.
I think you misunderstand.
v23 was built on v5, which is built on v1. Re-using the earlier logic was obviously better than duplicating it. v24 is used by an external system that nobody has any control over, so it’s impossible to change. All the other versions… well, no idea if anyone uses them, but everything works now, why invite disaster by removing any?
Does any developer here on HN really believe that JSON parsing (plus schema validation) is what adds the most latency to a request? It just doesn't add up that just switching to PB would deliver that speedup.
It reminds me of another headline I read a few years ago about a company thanking TypeScript for reducing bugs by 25% (I can't remember what the exact number was) after a rewrite... Common. You rewrote the whole thing from scratch. What actually cut the bugs by 25% was the rewrite itself. People don't usually write it more buggy the second time...
It's a pattern with these big tech consultants. They keep pushing specific tech solutions as silver bullets then they deliver flawed analyses of the results to propagate the myth that they solved the problem in that particular way... When in fact, the problem was solved by pure coincidence due to completely different factors.
OTOH if there was a lot of pressure to get the rewrite done, that would be conducive to producing buggy code. I think management would be a bigger factor than any technical issues.
I don't believe that for the purposes of inter-service communication, REST/JSON can compete with protobuf.
Caveat: I only have experience with REST/JSON and a little GraphQL, I haven't had the opportunity to work with protobuf yet. I'm more of a front end developer unfortunately, and I try to talk people out of doing microservices as much as possible.
For example, in JavaScript, using standard 'for loops' with index numbers is a LOT faster (over 16 times faster for basic use cases) than using Array.forEach(), yet all the linters and consultants recommend using Array.forEach() instead of the standard for loop... What about latency??? Suddenly nobody cares about latency nor performance.
The reason is that these operations which use marginal amounts of resources are pointless to refactor. If a function call which uses 1% of CPU time* (to service a standard request) has its performance improved by 'a whooping' 50%, then the program as a whole will use only 0.5% less CPU time than it did before.
* (where all the other operations use the remaining 99%)
Do you have proof of this claim? It smells like bs.
Because `Array.forEach` is one method call for each iteration, when `for loops` stay on the same frame. If the compiler can't inline that, it's a "major" overhead compared to simply jumping back at the top of the loop.
Googling for it I find some benchmark showing a 3x performance: https://leanylabs.com/blog/js-forEach-map-reduce-vs-for-for_..., but they call `Array.push` in the loop, so it's possible the difference is even bigger in practice.
To get an order of magnitude in difference, you'd have to construct a special case. I've seen a multiplier of about 3, but not 16.
As an aside: you use forEach() over for loops because most array processing operates on small arrays, where the marginal improvement of using a loop is limited. If you have a large array, you will eventually switch to a for loop if the loop body is relatively small. Likewise, when the request size is small, JSON works fine. But when your requests grows in size, the JSON overhead will eventually become a problem which needs a solution.
The underlying consideration is Amdahl's law.
Ideally you should compare running the loops without any logic inside them but I was worried that optimizations in the JavaScript engine would cause it to just skip over the loops if they performed no computations and without any memory side effects. Anyway this was beside the point I was trying to make. My point is already proven; it makes no sense to optimize cheap operations.
If your cache is hot and has a good hit-rate, the majority of your overhead is likely parsing. If you microbatch 100 requests, you have to parse 100 requests before you can ship them to the database for lookup (or the machine learning inference service). If the service is good at batch-processing, then the parsing becomes the latency-sensitive part.
Note the caveat: the 60% is for large payloads. JSON contains a lot of repetition in the data, so you often see people add compression to JSON unknowingly, because their webserver is doing it behind their back. A fairly small request on the wire deflates to a large request in-memory, and takes way more processing time.
That said, the statistician in me would like to have a distribution or interval rather than a number like "60%" because it is likely to vary. It's entirely possible that 60% is on the better end of what they are seeing (it's plausible in my book), but there's likely services where the improvement in latency is more mellow. If you want to reduce latency in a system, you should sample the distribution of processing latency. At least track the maximal latency over the last minute or so, preferably a couple of percentiles as well (95, 99, 99.9, ...).
From the article:
> The result of Protocol Buffers adoption was an average increase in throughput by 6.25% for responses and 1.77% for requests. The team also observed up to 60% latency reduction for large payloads.
It's very sneaky to describe throughput improvements using average requests/responses (which is what most people are interested in) but then switch to the 'worst case' request/response when describing latency... And doubly sneaky to then use that as the headline of the article.
There's also a lot of alarm bells going on when you have reports of averages without reports of medians (quartiles, percentiles) and variance. Or even better: some kind of analysis of the distribution. A lot of data will be closer to a Poisson-process or have multi-modality, and the average is generally hiding that detail.
What can happen is that you typically process requests around 10ms but you have a few outliers at 2500ms. Now the average is going to be somewhere between 10ms and 2500ms. If you have two modes, then the average can often end up in the middle of nowhere, say at 50ms. Yet you have 0 requests taking 50ms. They take either 10 or 2500.
In terms of formats, you'd get an easier transition and more balance between flexibility and efficiency out of BSON, Avro, Thrift, MessagePack. There are also alternatives to Protobuff like FlatBuffers and Cap'n Proto. There's also CBOR, which is interesting.
There are also other ways of looking at the problem. How does Erlang serialize messages? It doesn't because it messages itself, so the message format is native to itself. And in fact I mostly lean in that direction, but it's not for everyone. Erlang is also dynamically typed, not the kind of language Protobuff and Cap'n Proto is aimed at I suppose.
I don't get the difference you're drawing... the in-memory and on-the-wire representation of terms are different, so there's still serialization involved (term_to_binary/1). The format is documented and there are libraries for other languages.
Technically Erlang could go much further, but much like multicore support took forever, I guess due to lack of funding, it doesn't. Things like:
1. When transferring between two compatible nodes, or processes on the same node, 'virtually' serialize/deserialize skipping both operations and transferring pointer ownership to the other process instead.
2. When transferring between compatible nodes on different servers, use internal formats closer to mem representation rather than fully serializing to standard ETF/BERT
I've absolutely had times when json serialising/deserialising was the vast majority of the request time.
50k API endpoints means that probably a lot of them are pretty simple. The simpler the API, the higher percentage of it's time is spent in call overhead, which with JSON is parsing the entire input.
The article notes that this is only for "large payloads", likely an edge case, and the average performance improvement is 6.25% for responses and 1.77% for requests. 1.77%!
I get that this is "at scale" but is the additional complexity of all this engineering worth that? How much more difficult is the code to work with now? If it's at all more difficult to reason about, that is going to add to more engineering hours down the road to work with it.
I assume tradeoffs like this were taken into account, and it was deemed that a <7% response improvement was worth it.
Never tell them it's just because of all the technical debt you finally had organizational will to pay off.
To fix technical debt you often need a sexy sounding cover story.
What a plague of a company.
Some people will be fine finding jobs without it, a lot won't be that fortunate.
I’ve literally never seen a LinkedIn only job for a field I was interested in let alone a job - but I don’t work biz/fintech.
Sure, several of them, I could have found elsewhere, e.g. on a dedicated job board. But I wasn't at those other job boards, nor was the recruiter who brought them onto my radar.
If I had avoided LI, I would have missed out on some of my best gigs.
This is certainly wrong because the usage of LinkedIn is very high. It's equivalent to saying, "I don't take recruiting calls from meat-eaters." I mean, if it's extremely important to you, go for it, but if not, you're doing yourself a disservice.
You give... random applications... access to your contact list?
I think web integrations have gotten locked down since then.. not for altruistic reasons, but because Google and Apple don't want to give Microsoft access to the data they have on you, so yeah, it's at least harder to accidentally give LinkedIn access to your contacts list now, but the damage has already been done for a lot of people
And to come back to the current topic, their app felt so bloated that I removed it.
It's been a while since early Android, but I'm pretty sure that's how it worked.
A cursory glance over the article leads to some references to P99 latencies dropping from around mid 20ms to around mid 10ms.
Those are irrelevant.
Does the 1% care if their requests take 10ms more to fulfill?
From the surface this sounds like meaningless micro-optimization. I'm glad someone got to promote themselves with this bullshit though.
So feature wise, there's profile page of each user, there's searching page with filters, there's inbox, some recruitment tools I'd imagine, job posting, apply to the job and such along with search and recommendations for that too plus analytics for each of the component above certainly there would be some internal advertisement management platform.
What else I'm missing? So 50 thousand endpoints for all that?
So, do they not compress requests between microservices? Otherwise, how did they see such a reduction?
When I was at MS, it was definitely preferred over protobuf. IIRC it was also used internally as the wire format for gRPC. I guess LinkedIn is still kind of doing their own thing.
The others I have no idea about though.
But based on the complete lack of activity on the issue, gRPC + bond doesn’t seem like a popular choice. :)
Also, wouldn't something like MessagePack, ION, RION, or CBOR give you payload compression over JSON?
One thing about protobufs in a highly interconnected ball of mess, good luck reving the protocol in any non trivial way. Schema-less encodings such as those you mention (as well as bson, etc) are really advantageous for loosely coupled interfaces and graceful message format migrations. IMO (opinion!) protobufs is popular simply because of the Google cargo cult where anything Google tech is slavishly adopted without a great deal of introspection on applicability to the specific situation.
Duh! Of course going from a verbose text format like JSON to _any_ kind of binary format will lead to decreased everything.
Second instinct was:
Wait, were they using HTTP/2, HTTP/3? What JSON parser were they using? One of the fast ones, or one of the slow ones?
Third instinct was:
There is not enough information in the article to conclude "60%" of anything. It's not an apples to apples comparison. It's just an attention-seeking headline.
Then again, if it wasn't an attention seeking headline, it probably wouldn't have made it to HN. ;)
> Based on the learnings from the Protocol Buffers rollout, the team is planning to follow up with migration from Rest.li to gRPC, which also uses Protocol Buffers but additionally supports streaming and has a large community behind it.
At LinkedIn, we are focusing our efforts on advanced automation to enable a seamless, LinkedIn-wide migration from Rest.li to gRPC. gRPC will offer better performance, support for more programming languages, streaming, and a robust open source community. There is no active development at LinkedIn on new features for Rest.li. The repository will also be deprecated soon once we have migrated services to use gRPC. Refer to this blog[1] for more details on why we are moving to gRPC.
[0] - https://github.com/linkedin/rest.li
[1] - https://engineering.linkedin.com/blog/2023/linkedin-integrat...
I liked avro better though.
Don’t do this, it’s a terrible terrible idea.
It's crazy, but worked.
row database (bigtable), each team would have a column that stores their proto (or empty). So we'll run a batch process to read each such column, and then treat all the fields in the proto stored there as data, and expand all these fields as individual columns over the the columnar (read-only) db.
Later an analysts/statistician/linguist/etc. can query that columnar db with data they are interrested about. So that's what I remember (I left several years ago, so things might've changed), but pretty much instead of typical for row-databases have a column for everything - you just have a column for your protobuf (a bit like storing HSON/JSON in postgres), but then have the ETL process mow down through each field in that "column" and create "columns" for each such field.
We had to do some workarounds though as it was exporting too much, and it was not clear how we can know which fields would be asked about (actually there was a way probably, but it'll taken time to coordinate with some other team, so it was more (I think) on coordination with internal "customers" to disallow exporting these).
But the cool thing, is that if customer team added new field to their proto, our process would see it, and it'll get expanded. If they deprecate a proto, there could be (not sure if there was, but could be added) - no longer export it. But for this to work you need the "protodb" e.g. to introspect, able to reflect actual names in order to generate the column.
But that's exactly the problem with protocol buffers. In typical Google fashion, they never bothered to create support for the native language of their platform: Kotlin. And that was just the one example we encountered on a project. If your language isn't supported by protoc, then what?
Also touted was protocol buffers' "future-proofing" or resistance to version changes. We never realized this alleged benefit, and it wasn't clear how it was supposed to be delivered.
I dreaded the opacity of the format, and expected decodes to start failing in mysterious ways at some point... wasting days or weeks of my time. But I have to admit that this rarely if ever happened. In our case we were using C++ on one end and Swift on the other.
I don't get this rationale in the article: "While the size can be optimized using standard compression algorithms like gzip, compression and decompression consumes additional hardware resources." And unpacking protobufs doesn't?
I might define a different class for a validated object, so version skew only needs to be dealt with at the system edge.
Maybe you don’t actually have that problem because you can make a closed-world assumption? For example, you control all writers and aren’t saving any old-format data, like in a database that has a schema. In that case you can make fields required.
https://cloud.google.com/blog/products/application-developme...
https://developers.googleblog.com/2021/11/announcing-kotlin-...
1) Never change the type or semantic meaning of a field.
It's really that simple.
"Reason #2: Backward Compatibility For Free
Numbered fields in proto definitions obviate the need for version checks which is one of the explicitly stated motivations for the design and implementation of Protocol Buffers... With numbered fields, you never have to change the behavior of code going forward to maintain backward compatibility with older versions."