How to force Zookeeper server to wait until client reads response?

Viewed 358

I have a Zookeeper cluster created of several nodes running on kubernetes. I have also one "manager" pod which should fetch Zookeeper nodes metadata. I am achieving this by sending a stat word to the zookeeper nodes. After sending a request I should be able to read response with metadata but unfortunately, this is not true all the time. There is some kind of race condition because sometimes reading response is performed without IOException exception, sometimes the exception is thrown in the middle of reading the data (the data are received but the exception is thrown nevertheless). And there is the third option - no data are received and the exception is thrown.

I think the Zookeeper node closes a connection before/in the middle of reading response. Is there any way how to force Zookeeper to wait until the response is fully received?

private void getMetadataFromPod(Pod pod, SSLSocketFactory factory) {
int port = 2181;
    try {
        String host = getPodIP();
        SSLSocket socket = null;
        BufferedReader in = null;
        PrintWriter out = null;
        SSLSocket socket = (SSLSocket) factory.createSocket();

        if (socket == null) {
            log.error("Could not create socket for getting Zookeeper data");
            return;
        }
        try {
            socket.connect(new InetSocketAddress(host, port), 10_000);
        } catch (ConnectException | SocketTimeoutException e) {
            log.error("Could not connect " + e.getMessage());
        }
        try {
            log.debug("Starting handshake with {}", socket.getRemoteSocketAddress());
            try {
                socket.startHandshake();
                out = new PrintWriter(
                        new BufferedWriter(
                                new OutputStreamWriter(
                                        socket.getOutputStream(), StandardCharsets.UTF_8)));
                out.println("stat");
                out.flush();

                in = new BufferedReader(
                        new InputStreamReader(
                                socket.getInputStream(), StandardCharsets.UTF_8));
                String inputLine;
                while ((inputLine = in.readLine()) != null) {
                    log.debug(inputLine);
                }
            } catch (SSLHandshakeException e) {
                log.error("Error while performing TLS handshake with pod {} in namespace {}",
                        pod.getMetadata().getName(), pod.getMetadata().getNamespace(), e);
            } catch (IOException e) {
                e.printStackTrace();
            }
        } finally {
            log.debug("Closing resources");
            in.close();
            out.close();
            socket.close();
        }
    } catch (IOException e) {
        log.debug("Error while getting Zookeeper metadata: " + e.getMessage());
    }
}

Zookeeper server log from the case where the exception on the client side was thrown

2019-01-11 10:55:33,907 INFO Accepted socket connection from /127.0.0.1:46698 (org.apache.zookeeper.server.NIOServerCnxnFactory) [NIOServerCxn.Factory:0.0.0.0/0.0.0.0:21813]
2019-01-11 10:55:33,908 INFO Processing stat command from /127.0.0.1:46698 (org.apache.zookeeper.server.NIOServerCnxn) [NIOServerCxn.Factory:0.0.0.0/0.0.0.0:21813]
2019-01-11 10:55:33,909 INFO Stat command output (org.apache.zookeeper.server.NIOServerCnxn) [Thread-811]
2019-01-11 10:55:33,910 INFO Closed socket connection for client /127.0.0.1:46698 (no session established for client) (org.apache.zookeeper.server.NIOServerCnxn) [Thread-811]

Log from "manager" client

java.net.SocketException: Connection reset
at java.net.SocketInputStream.read(SocketInputStream.java:210)
at java.net.SocketInputStream.read(SocketInputStream.java:141)
at sun.security.ssl.InputRecord.readFully(InputRecord.java:465)
at sun.security.ssl.InputRecord.read(InputRecord.java:503)
at sun.security.ssl.SSLSocketImpl.readRecord(SSLSocketImpl.java:975)
at sun.security.ssl.SSLSocketImpl.readDataRecord(SSLSocketImpl.java:933)
at sun.security.ssl.AppInputStream.read(AppInputStream.java:105)
at sun.nio.cs.StreamDecoder.readBytes(StreamDecoder.java:284)
at sun.nio.cs.StreamDecoder.implRead(StreamDecoder.java:326)
at sun.nio.cs.StreamDecoder.read(StreamDecoder.java:178)
at java.io.InputStreamReader.read(InputStreamReader.java:184)
at java.io.BufferedReader.fill(BufferedReader.java:161)
at java.io.BufferedReader.readLine(BufferedReader.java:324)
at java.io.BufferedReader.readLine(BufferedReader.java:389)
at myproject.getMetadataFromPod(MyClass.java:295)
at myproject.MyClass.lambda$zookeeperData$5(MyClass.java:337)
at io.vertx.core.impl.ContextImpl.lambda$executeBlocking$1(ContextImpl.java:273)
at io.vertx.core.impl.TaskQueue.run(TaskQueue.java:76)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
at java.lang.Thread.run(Thread.java:748)
0 Answers
Related