RPC - Remote Procedure Call
Used to solve communication in distributed systems. The core idea is that you can make a remote call the same way you call something local. RPC isn't just a microservice / cloud-native term — anywhere there's network communication, you might use RPC.
Two examples:
- A large distributed app might depend on a message queue, a distributed cache, a distributed database, a unified config center, and so on. The app can talk to those middlewares over RPC. etcd, for example, as a unified config service: the client talks to the server through the gRPC framework.
- Kubernetes itself is distributed. Communication between kube-apiserver and every component in the cluster goes through gRPC.
RPC involves:
- Serialization: turn objects into a transmittable byte stream (serialize) and reverse it (deserialize). Solves data exchange across the network and across languages.
- Compression: shrink what's sent over the wire, cut bandwidth and latency.
- Protocol: rules for client–server communication, including wire format and interaction mode: HTTP/2, TCP, UDP.
- Dynamic proxy: hide the complexity of remote calls so developers use a remote service like a local method: JDK dynamic proxy, bytecode enhancement.
- Service registration and discovery: dynamically manage whether instances are available, support load balancing and failover. Registries like ZooKeeper, Consul, ETCD store addresses and metadata.
- Encryption: confidentiality and integrity of data in transit; stop MITM and tampering.
- Network I/O: I/O models for efficient, stable communication — connection management, send/receive, the low-level stuff. Sounds simple; actually it's messy: finding the peer, establishing a connection, encoding/decoding, managing connections. RPC wraps that whole path, so when you build a distributed system the networking logic is simpler, and also more secure and reliable.
RPC clusters also involve:
- Monitoring
- Circuit breaking / rate limiting
- Graceful start and stop
- Multiple protocols
- Distributed tracing
Where RPC is actually powerful:
- Connection management
- Health checks
- Load balancing
- Graceful startup and shutdown
- Retry on failure
- Business grouping
- Circuit breaking / rate limiting
Without an RPC framework, how would you even call an API on another machine?
RPC hides the network-programming details so calling a remote method feels like calling a local one (a method in the same project). You don't have to write a pile of non-business code just because the method is remote.
RPC mainly does two things:
- Hide the difference between remote and local calls, so it feels like a method in this project;
- Hide how messy the underlying network is, so you can stay on business logic.
Serialization
Data on the wire has to be binary, but the caller’s in/out args are objects. You convert them to transmittable binary first, and the conversion has to be reversible.
The header is usually for identity: protocol id, size, request type, serialization type, etc. The body is mainly business params and extra attributes.
Deserialization


RPC doesn't only solve communication — you can also use it toward MQ, distributed cache, databases.
RPC and HTTP are both application-layer protocols.
Before an RPC request goes on the network, it turns the method-call args into binary; then writes that into a local Socket, and the NIC sends it out.

For an extensible, backward-compatible protocol, the key is using the extension fields in the Header and in the Payload, and staying compatible through those fields.
Pick a serialization method that fits the scenario.

Common serialization methods:
- JDK native serialization

Any serialization framework, at heart, is designing a serialization protocol.
- JSON: typical key-value, no types, a text serialization framework.
Two problems with JSON serialization:
- Extra space cost is high — for large payloads that means huge memory and disk;
- JSON has no types, so a strongly typed language like Java has to go through reflection, which is slow.
So if an RPC framework picks JSON, the data between provider and caller should stay relatively small, or performance tanks.
- Hessian: dynamically typed, binary, compact, portable across languages. More compact than JDK and JSON, much faster, and fewer bytes.
- Hessian itself has issues. The official version doesn't support some common Java types.
- Linked* — LinkedHashMap, LinkedHashSet, etc. You can fix this by extending CollectionDeserializer;
- Locale — extend ContextSerializerFactory;
- Byte/Short become Integer on deserialize
- Hessian itself has issues. The official version doesn't support some common Java types.
- Protobuf: Google’s internal mixed-language data standard, a structured storage format, usable for serializing structured data. Supports Java, Python, C++, Go, etc. You define an IDL (Interface description language), then use per-language IDL compilers to generate helpers. Pros:
- Much smaller than JSON / Hessian after serialize;
- IDL describes semantics clearly, so types don't get lost between apps — no XML-style parser;
- Serialize/deserialize is fast; no reflection for types;
- Message format upgrades and compatibility are decent; can be backward compatible.
There's a Java-oriented Protobuf-like framework that doesn't need an IDL file and can deserialize Java domain objects directly. Efficiency is about the same as Protobuf, binary format is identical — basically a Java Protobuf. In practice I hit some unsupported cases:
Other protocols: Message Pack, kryo, etc.
What affects the choice:

Default pick is still Hessian and Protobuf — they cover performance, time, space, generality, compatibility, and safety. Hessian is easier to use and better on object compatibility; Protobuf is more efficient and more general.
What to watch when using an RPC framework?
- Objects built too nested / too many fields.
- Objects too huge.
- Using a class the serializer doesn't support as an input type.
- Complicated inheritance.
Which network I/O model does RPC tend to use?
Common I/O models
- Blocking I/O (BIO)
- Non-blocking I/O (NIO)
- I/O multiplexing
- Asynchronous non-blocking I/O (AIO)
Only AIO is async I/O; the rest are sync.
Blocking I/O is the simplest, most common model. On Linux, sockets are blocking by default. Flow:
The process makes an I/O syscall and blocks, work moves to kernel space. The kernel waits for data, then copies it into user memory, then I/O is done and returns. Then the process unblocks and runs business logic.
Kernel I/O has two phases — wait for data, copy data. In both, the I/O thread in the app stays blocked. If you're on Java threads, every I/O holds a thread until it finishes.
I/O multiplexing
The most widely used model under high concurrency. Java NIO, Redis, Nginx underneath are this; classic Reactor is based on it too.
I/O from many connections can register on one multiplexer (select). When the user process calls select, the whole process blocks. The kernel “watches” all those sockets; when any socket has data ready, select returns. Then the user process calls read and copies from kernel to user space.
So: select blocks until some socket is ready, then you read. The path is more complicated than blocking I/O, looks like more overhead. The win is you can handle many sockets’ I/O in one thread. Register many sockets, keep calling select, read the ones that fired — same thread, many I/O requests. Blocking I/O needs many threads for that.
Why are blocking I/O and I/O multiplexing the usual ones?
You need kernel support and language support.
Most kernels support blocking I/O, non-blocking I/O, and multiplexing. Signal-driven I/O and async I/O only show up on newer Linux.
In C++ and Java, high-performance network frameworks are mostly Reactor, Netty being the typical Java one. Reactor is multiplexing. In non-hot paths, blocking I/O is still the most common.
Which I/O model for RPC?
Most RPC is high-concurrency. Given kernel support, language support, and the models themselves, RPC implementations pick I/O multiplexing. For the language’s network framework, best is a Reactor implementation — in Java, Netty (there are other NIO frameworks; Netty is the one people actually use). On Linux, turn on epoll (Windows can't; the kernel doesn't have it).
What is a Reactor-based network I/O model?
A Reactor I/O model is an event-driven high-performance networking model. It decouples listening for I/O events, dispatching them, and business logic, so many connections can be managed and answered in one place. Core: multiplexing (Select, epoll, kqueue) watches many connection events, and dispatches by event type to a handler — no thread-per-connection waste of blocking I/O.
Core pieces:
- Reactor: listens for all I/O events, and in an Event Loop dispatches ready events to handlers. Hub of the model, usually its own thread. Uses a multiplexer (Selector) to poll registered Channels for connect/read/write.
- Acceptor: handles new connections, accepts the client, registers the new SocketChannel on the Reactor for later read/write.
- Handler: actual work (read, process, write back), typically:
Read handler: read-ready, read from Channel and decode. Write handler: write-ready, encode the result and write back. Business handler: compute, DB, slow stuff, maybe a thread pool.
Zero copy
Kernel I/O: wait for data, copy data. Waiting: NIC gets data, kernel writes it in. Copying: kernel copies into the user process.

Every write from the app goes to a user-space buffer, CPU copies to kernel buffer, DMA copies to the NIC, NIC sends. Two copies before it leaves. Reads are the reverse — two copies again before the app sees it.
A full read/write bounces between user and kernel, and every copy is a CPU context switch (user → kernel or kernel → user).
Zero-copy
Zero-copy means dropping the copy between user space and kernel space. Each read/write is done so writing/reading user space is like writing/reading kernel space, then DMA copies kernel ↔ NIC.

Two approaches
- mmap+write: virtual memory.
- sendfile
mmap+write
How:
- mmap maps the kernel read buffer into the process’s virtual address space, shared memory. No copy kernel → user buffer, just an address mapping.
Transfer:
- First copy (DMA): disk → kernel read buffer.
Shared mapping: user accesses kernel buffer via the mapping.
- Second copy (CPU): on write, CPU copies kernel read buffer → kernel Socket buffer.
- Third copy (DMA): Socket buffer → NIC.
Pros/cons:
Pros:
- One less CPU copy (kernel → user buffer is gone).
- App can touch the mapped memory; good if you need to preprocess (edit, compress).
Cons:
- Still 4 context switches (two syscalls) and 3 copies.
- Mapping has overhead; truncating the file can blow up (SIGBUS).
sendfile
sendfile merges read and write into one syscall; transfer stays in kernel.
Transfer (two modes)
Basic (no SG-DMA):
- First copy (DMA): disk → kernel read buffer.
- Second copy (CPU): kernel read buffer → kernel Socket buffer.
- Third copy (DMA): Socket buffer → NIC.
SG-DMA:
Only two DMA copies: kernel read buffer goes to the NIC via DMA Scatter/Gather; CPU doesn't copy into the Socket buffer.
Pros:
- Syscalls down to 1, context switches 2.
- With SG-DMA, real zero-copy (two DMA copies only).
- Throughput jumps; good for big files.
Cons:
- Data is invisible to user space; you can't process it before sending.
- Needs OS and hardware support.
If you need to preprocess (edit the file), mmap+write.
If you just want fast transfer and don't touch the data, prefer sendfile (especially with SG-DMA).
Zero-copy in Netty
This is entirely in user space, i.e. the JVM. Netty’s “zero-copy” is mostly about optimizing how you operate on data.
- CompositeByteBuf: merge several ByteBufs into one logical ByteBuf, no copy between them.
- ByteBuf slice: split a ByteBuf into several that share the same storage, no copy.
- wrap: wrap byte[], ByteBuf, ByteBuffer as a Netty ByteBuf, no copy.
Netty also wraps NIO FileChannel.transferTo() in FileRegion — same idea as Linux sendfile.
Dynamic proxy: program to interfaces, hide the RPC pipeline (I haven't actually read the code, so I'm not that clear on this)
On networking, just remember — reliable transport.
RPC auto-generates a proxy for the interface. When you inject the interface, at runtime what's bound is that generated proxy. Calls get intercepted by the proxy, and that's where the remote-call logic lives.

- Proxies are generated at runtime, so generation speed, bytecode size, etc. all hit performance — smaller bytecode, less runtime cost.
- The proxy intercepts every interface call, so it has to run fast.
- You want a proxy framework that's easy to use. API, community, dependency complexity.
gRPC

Protocol framing
After the binary of the method args you need “sentence breaks” to split requests. What's between two breaks is one request's binary. That's protocol encapsulation.


Service discovery: CP or AP?

- Register: when a provider starts, it registers the exposed interfaces in the registry; the registry keeps that node's IP and interfaces.
- Subscribe: when a caller starts, it looks up and subscribes to provider IPs, caches them locally, uses them for later remote calls.

If you use DNS for discovery:
All provider nodes under one domain, caller can get a random provider IP via DNS and hold a long connection. Looks fine, until:
- If an IP:port goes down, can the caller drop that node in time?
- If you already had some nodes up and then scale out, can the new nodes get traffic in time?
Answer is no. For performance and to spare DNS, DNS is heavily cached, usually for a long time.
ZooKeeper-based discovery

- The platform admin creates a service root in ZooKeeper, often named after the interface (e.g.
/service/com.demo.xxService), then provider and consumer dirs under that, to store provider and caller node info. - When a provider registers, it creates an ephemeral node under the provider dir with its registration info.
- When a caller subscribes, it creates an ephemeral node under the consumer dir with its info, and watches all service nodes under the provider dir (
/service/com.demo.xxService/provider). - When data under the provider dir changes, ZooKeeper notifies subscribed callers.
An eventually-consistent registry on a message bus
ZooKeeper’s big trait is strong consistency. Every update on a node is applied on the others at the same time. Real-time identical data on every node — that's why ZooKeeper cluster performance drops.
For RPC discovery, when a node just came up, callers can live with noticing it a few seconds later. A few seconds (or more) without traffic after a node starts doesn't matter to the cluster. So we can drop CP (strong consistency) and take AP (eventual consistency) for registry performance and stability.
If you want eventual consistency, a message bus is an option. Registry data can be fully cached in each registry process, synced over the bus. One node receives a registration, publishes to the bus, others update and push to callers — eventual consistency between registries:

Afterword
Later I should walk through some gRPC and Kitex code, and there are ByteDance cloud-native WeChat posts worth reading. The actually important thing is to put the open-source projects aside and get the fundamentals.