Сервис перенаправляющий grpc запросы

В настоящее время я работаю над сервером grpc, который будет получать потоковые вызовы grpc с первого сервера и перенаправлять эти вызовы на второй сервер, а также перенаправлять ответы со второго сервера в виде потоков на первый.

У меня есть 2 прото файла, первый прото

Первый файл:

syntax = "proto3";

package first.proto.pack;

service FirstProtoService {
  rpc StreamingCall(stream RequestToFirstServer) returns (stream ResponseForFirstServer){}
}

message RequestToFirstServer {
    oneof firstStreamingRequest {
        int32 x = 1;
        int32 y = 2;
    }
}

message ResponseForFirstServer {
  string someprocessedinformation = 1;
}

Второй файл:

syntax = "proto3";

package second.proto.pack;

service SecondProtoService {
  rpc StreamingCall(stream RequestToSecondServer) returns (stream ResponseFromSecondServer){}
}

message RequestToSecondServer {
  oneof secondStreamingRequest {
    int32 processedX = 1;
    int32 procesdedY = 2;
  }
}

message ResponseFromSecondServer {
  string computedInformation = 1;
}

Первый сервер знает о первом прото-файле, но не знает о втором.

Второй сервер знает о втором прото-файле, но не знает о первом.

Средний сервер знает о первом и втором протоколе.

Необходимо написать сервер, который будет передавать запросы с одного сервера с одного сервера на другой.

Я начал писать на Java. Но столкнулся с проблемой отправки большого количества запросов на второй сервер

Вот как моя реализация среднего уровня обслуживания выглядит на Java:

package middle.server.pack;

import first.proto.pack.First;
import first.proto.pack.FirstProtoServiceGrpc;
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.stub.StreamObserver;
import second.proto.pack.Second;
import second.proto.pack.SecondProtoServiceGrpc;

import java.util.logging.LogManager;
import java.util.logging.Logger;

public class MiddleService extends FirstProtoServiceGrpc.FirstProtoServiceImplBase {
    private final ManagedChannel channel = ManagedChannelBuilder.forTarget("localhost:8080").build();
    private final Logger logger = LogManager.getLogManager().getLogger(MiddleService.class.getName());

    @Override
    public StreamObserver<First.RequestToFirstServer> streamingCall(StreamObserver<First.ResponseForFirstServer> responseObserver) {
        return new StreamObserver<First.RequestToFirstServer>() {
            @Override
            public void onNext(First.RequestToFirstServer value) {
                SecondProtoServiceGrpc.SecondProtoServiceStub stub = SecondProtoServiceGrpc.newStub(channel);
                StreamObserver<Second.RequestToSecondServer> requestObserver = stub.streamingCall(
                        new StreamObserver<Second.ResponseFromSecondServer>() {
                            @Override
                            public void onNext(Second.ResponseFromSecondServer value) {
                                doProcessOnResponse(value);
                                First.ResponseForFirstServer responseForFirstServer =
                                        mapToFirstResponse(value);
                                responseObserver.onNext(responseForFirstServer);
                            }

                            @Override
                            public void onError(Throwable t) {
                                logger.info(t.getMessage());
                            }

                            @Override
                            public void onCompleted() {
                                logger.info("sucess");
                            }
                        }
                );
                Second.RequestToSecondServer requestToSecondServer = mapToSecondRequest(value);
                requestObserver.onNext(requestToSecondServer);
                requestObserver.onCompleted();
            }

            @Override
            public void onError(Throwable t) {
                logger.info(t.getMessage());
            }

            @Override
            public void onCompleted() {
                logger.info("Everything okay");
            }
        };
    }
}

После запроса с первого клиента на стороне серединного сервера получаю такие ошибки:

CANCELLED: Failed to read message.
CANCELLED: io.grpc.Context was cancelled without error

Я знаю, что делаю это неправильно. Итак, вопрос в том, как это исправить, или, если я не могу сделать это на java, могу ли я сделать это на любом другом языке?


Ответы (0 шт):