Zero copy in Java & Go
- How to send file in java
- How to send file in go
- Java nio zero copy file transfer
- Golang zero copy with io.Copy
- Java FileChannel transferTo timeout
- Grpc streaming vs zero copy for large files
- Reliable large file transfer with Java nio selectors
- Architecture for parallel large file transfer
Suppose that we need to transfer a multi-gigabyte/terabyte file between two peers (Java/Go apps) over the network.
What ideas come to mind at first? The first thoughts would be to:
- Split file into chunks and transfer them using HTTP
- Use gRPC streaming (HTTP/2), adjust payload size
gRPC’s default payload size is 4 mb, and it is a security constraint, not the protocol limit – it can be adjusted on the server and client side.
For some cases, gRPC is a “good enough” solution. It has predicatable behavior and usually takes the most part of the work, related to reliability.
However, such an approach has downsides:
- For the maximum performance – payload size should be adjusted in according to system capabilities
Also, it puts those batches to the heap/direct memory of the application, prior to sending/receiving them:
- In the case of multiple parallel transfers (suppose, 1000), and certain payload size (suppose 80 MB), we would have to allocate min 80 GB of RAM!
- Since ultimately, we will have to pass everything through the memory of the app – we waste a lot of resources: data needs to be copied to the app memory, then GC needs to remove it, and all over again.
The trick here is to employ zero-copy mechanism – we avoid passing the file through the app memory, making OS kernel to tackle this tranfer directly. Moreover, we don’t even need to split a file into chunks, or somehow care about payload!
We just start streaming the file, and OS kernel and TCP stack handle the rest: backpressure, ordering, and maybe even brew some coffee while you lean back and relax.
General architecture
Suppose that Bob wants to receive a directory with files from Alice as fast as possible. Bob doesn’t know the structure, but he is sure that there may be hundreds of 1 TB files.
Here is Bob’s plan:
- Bob sends gRPC request to Alice
listDirectory(dir_name), asking for a directory structure: file list, with names and sizes. - Bob receives that list and spins up 100 parallel threads.
- In each thread, bob opens a socket on a random free port.
- Each thread sends gRPC request to Alice
sendFile(file_name, ip, port), letting Alice know that Bob is ready to receive a certain file on a certain port. - When Alice receives such a gRPC request, she initiates streaming in according to given input.
- Neither Alice nor Bob uses application memory; OS kernel dynamically adjusts the stream intensity and payload size.
Java
Java has IO and NIO. It’s not possible to make zero-copy file transfer with IO, so we will be using NIO, but there is a catch.
Simple file receiver
try (ServerSocketChannel serverChannel = ServerSocketChannel.open()) {
serverChannel.socket().bind(new InetSocketAddress(port));
File outputFile = new Package.File(outputFilePath);
try (SocketChannel clientChannel = serverChannel.accept();
FileOutputStream fos = new FileOutputStream(outputFile);
FileChannel destChannel = fos.getChannel();) {
long bytesTransferred = 0;
while (bytesTransferred < expectedFileSize) {
// Transfer bytes from the client's socket channel into our file channel directly
long transferred = destChannel.transferFrom(clientChannel, bytesTransferred, expectedFileSize - bytesTransferred);
if (transferred <= 0) {
throw new IOException("connection closed unexpectedly");
}
bytesTransferred += transferred;
}
}
}
Simple file sender
long fileSize = sourceFile.length();
try (SocketChannel socketChannel = SocketChannel.open(new InetSocketAddress(host, port));
FileInputStream fis = new FileInputStream(sourceFile);
FileChannel sourceChannel = fis.getChannel()) {
long bytesSent = 0;
while (bytesSent < fileSize) {
// Transfer bytes from our file channel into the socket channel directly
bytesSent += sourceChannel.transferTo(bytesSent, fileSize - bytesSent, socketChannel);
}
}
This is a working POC – you can run it and see it working.
But now you may ask questions – what if the receiver opens a socket and waits for the sender, but the sender never arrives? What if any of the parties dies midway? And those are the goddamn right questions: the process may end up in a deadlock.
And it turns out that implementing timeout – to make this system reliable – is a whole different story. Let’s reiterate, what types of socket timeouts we have:
- Socket connection timeout: the receiver opens a socket, and waits sender to connect
- Socket read timeout: the sender is connected, but doesn’t send data in time.
It turns out that Java’s transferTo() doesn’t natively support timeouts.
The trick here is to use:
- IO capabilities for the connection establishment, thereby allowing connection timeouts to be configured.
- Use non-blocking NIO selectors to be able to set read timeout.
Without further ado, here is receiver’s and sender’s code for reliable zero-copy transfer.
Reliable file receiver
try (ServerSocketChannel serverChannel = ServerSocketChannel.open()) {
ServerSocket serverSocket = serverChannel.socket();
serverSocket.setSoTimeout(CONNECT_TIMEOUT_MS);
serverSocket.bind(new InetSocketAddress(port));
try (SocketChannel clientChannel = serverSocket.accept().getChannel();
Selector selector = Selector.open();
FileOutputStream fos = new FileOutputStream(outputFilePath);
FileChannel destChannel = fos.getChannel()
) {
clientChannel.configureBlocking(false);
clientChannel.register(selector, SelectionKey.OP_READ);
long bytesTransferred = 0;
long lastBatchStampMs = System.currentTimeMillis();
while (bytesTransferred < expectedFileSize) {
int readyChannels = selector.select(READ_TIMEOUT_MS);
if (readyChannels == 0) {
// The select() call timed out; no data arrived in time
throw new SocketTimeoutException("Transfer stalled: no data received for " + READ_TIMEOUT_MS + " ms");
}
selector.selectedKeys().clear();
long transferred = destChannel.transferFrom(clientChannel, bytesTransferred, expectedFileSize - bytesTransferred);
if (transferred > 0) {
bytesTransferred += transferred;
lastBatchStampMs = System.currentTimeMillis();
} else if (transferred == -1) {
throw new IOException("connection closed prematurely by sender");
} else if (transferred == 0) {
if (System.currentTimeMillis() - lastBatchStampMs > READ_TIMEOUT_MS) {
throw new IOException("no progress for a long time");
}
} else {
throw new Exception("unexpected bytes value: " + transferred);
}
}
}
}
Note that transferFrom() may return value will be zero in some cases. For example, when sender’s resources are exhausted, network congestion occurs, etc.
But it is also possible that the sender is dead. Therefore, we need to track the duration during which no data is received and close the connection ourselves if the timeout is exceeded.
Reliable file sender
long fileSize = sourceFile.length();
try (SocketChannel socketChannel = SocketChannel.open()) {
Socket socket = socketChannel.socket();
socket.connect(new InetSocketAddress(host, port), CONNECT_TIMEOUT_MS);
try (FileInputStream fis = new FileInputStream(sourceFile);
FileChannel sourceChannel = fis.getChannel()
) {
long bytesSent = 0;
while (bytesSent < fileSize) {
// One may want to set timeout here as well
bytesSent += sourceChannel.transferTo(bytesSent, fileSize - bytesSent, socketChannel);
}
}
}
Under the hood, on Linux and other OS, this method can leverage the sendfile syscall.
Golang
In golang we have the io.Copy function. It automatically chooses the most efficient data transfer mechanism available.
When io.Copy detects that it’s copying from a file to a TCP connection, it will automatically attempt to use zero-copy syscalls like sendfile or splice on Linux, avoiding that need to pass data through the application’s user-space memory.
This means we can achieve our goal without manual chunking or buffer management.
We use net.DialTimeout to establish a connection with a connect timeout, just like the Java example. Then, a single call to io.Copy handles the entire file transfer for zero-copy performance.
The sender doesn’t need to implement a write timeout loop; if the receiver times out and closes the connection, the sender’s io.Copy will fail with a “broken pipe” error, preventing it from blocking indefinitely.
Reliable file receiver
listener, err := net.Listen("tcp", ":"+port)
if err != nil {
return err
}
defer listener.Close()
conn, err := listener.Accept()
if err != nil {
return err
}
defer conn.Close()
outputFile, err := os.Create(outputFilePath)
if err != nil {
return err
}
defer outputFile.Close()
// pay attention – it is a deadline, not the timeout
if err := conn.SetReadDeadline(time.Now().Add(READ_TIMEOUT_RECEIVER)); err != nil {
return err
}
_, err = io.Copy(outputFile, conn)
if err != nil {
if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
return errors.New("transfer stalled: timeout exceeded")
}
return err
}
Reliable file sender
sourceFile, err := os.Open(sourceFilePath)
if err != nil {
return err
}
defer sourceFile.Close()
conn, err := net.DialTimeout("tcp", net.JoinHostPort(host, port), CONNECT_TIMEOUT_SENDER)
if err != nil {
return err
}
defer conn.Close()
// pay attention – it is a deadline, not the timeout
if err := conn.SetWriteDeadline(time.Now().Add(WRITE_TIMEOUT_SENDER)); err != nil {
return err
}
if _, err := io.Copy(conn, sourceFile); err != nil {
return err
}
That’s all.