gRPC ja suorituskykyiset API:t · Oppitunti

Asiakaspohjainen streamaus

Opi toteuttamaan asiakaspään streamaus, jossa asiakas voi lähettää palvelimelle viestisarjan.

Oppitunti 2/410 vaihetta

Asiakaspohjainen streamaus on ilmainen gRPC ja suorituskykyiset API:t-oppitunti CoddyKitissä. Tämä on oppitunti 2/4. Voit lukea koko oppitunnin alta ilmaiseksi ja harjoitella sen jälkeen käytännössä selaimessa sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla. Oppitunti kuuluu gRPC ja suorituskykyiset API:t-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. gRPC ja suorituskykyiset API:t-kurssilla on yhteensä 4 oppituntia.

Mitä asiakaspuolen suoratoisto on?

Tervetuloa tutustumaan asiakaspuolen suoratoistoon! gRPC:ssä asiakaspuolen suoratoisto on viestintämalli, jossa asiakas lähettää palvelimelle viestisarjan.

Yksinkertaisessa unary-RPC:ssä (pyyntö–vastaus) asiakas ei lähetä vain yhtä viestiä. Sen sijaan se lähettää viestivirran, jonka palvelin käsittelee ja lähettää lopuksi yhden vastauksen.

Miten asiakaspuolen suoratoisto toimii

Kuvitelkaa, että lataatte suuren tiedoston lähettämällä sen monina pieninä osina. Palvelin kerää kaikki osat, kokoaa tiedoston uudelleen ja lähettää sitten yhden "lataus valmis" -vahvistuksen.

  • Asiakas aloittaa RPC-kutsun.
  • Asiakas lähettää useita viestejä asynkronisesti.
  • Palvelin vastaanottaa ja käsittelee nämä viestit.
  • Kun asiakas lopettaa lähettämisen (ilmoittaa suorittamisen päättymisestä), palvelin lähettää yhden vastauksen.

Asiakaspuolen suoratoiston määrittäminen Protobufissa

Voitte määrittää asiakaspuolen suoratoistometodin .proto-tiedostossa käyttämällä stream-avainsanaa request-tyypille, mutta ette vastaustyypille.

Tässä on esimerkki lokien latauspalvelusta:

syntax = "proto3";

package client_streaming;

service LogService {
  rpc UploadLogs (stream LogEntry) returns (UploadSummary);
}

message LogEntry {
  string message = 1;
  int64 timestamp = 2;
}

message UploadSummary {
  int32 uploaded_count = 1;
  string status_message = 2;
}

Palvelin: asiakaspuolen suoratoiston käsittely

Palvelinpuolella metodinne vastaanottaa StreamObserver-olion asiakkaalta saapuvia viestejä varten ja käyttää toista StreamObserver-oliota yksittäisen vastauksen lähettämiseen.

Palvelimen StreamObserver-oliolla on onNext()-metodi jokaista saapuvaa viestiä varten, onError()-metodi virheitä varten ja onCompleted()-metodi, kun asiakas lopettaa lähettämisen.

import io.grpc.stub.StreamObserver;
import io.grpc.Server;
import io.grpc.ServerBuilder;
import client_streaming.LogEntry;
import client_streaming.LogServiceGrpc;
import client_streaming.UploadSummary;

public class LogServer {
    private Server server;

    private void start() throws Exception {
        int port = 50051;
        server = ServerBuilder.forPort(port)
            .addService(new LogServiceImpl())
            .build()
            .start();
        System.out.println("Server started, listening on " + port);
        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            System.err.println("*** shutting down gRPC server since JVM is shutting down");
            LogServer.this.stop();
            System.err.println("*** server shut down");
        }));
    }

    private void stop() {
        if (server != null) {
            server.shutdown();
        }
    }

    private void blockUntilShutdown() throws InterruptedException {
        if (server != null) {
            server.awaitTermination();
        }
    }

    public static void main(String[] args) throws Exception {
        final LogServer logServer = new LogServer();
        logServer.start();
        logServer.blockUntilShutdown();
    }

    static class LogServiceImpl extends LogServiceGrpc.LogServiceImplBase {
        @Override
        public StreamObserver<LogEntry> uploadLogs(StreamObserver<UploadSummary> responseObserver) {
            return new StreamObserver<LogEntry>() {
                private int logCount = 0;

                @Override
                public void onNext(LogEntry log) {
                    // Process each log entry as it arrives
                    System.out.println("Received log: " + log.getMessage() + " at " + log.getTimestamp());
                    logCount++;
                }

                @Override
                public void onError(Throwable t) {
                    System.err.println("UploadLogs cancelled or failed: " + t.getMessage());
                    responseObserver.onError(t);
                }

                @Override
                public void onCompleted() {
                    // After all logs are received, send a single summary response
                    UploadSummary summary = UploadSummary.newBuilder()
                        .setUploadedCount(logCount)
                        .setStatusMessage("Successfully processed " + logCount + " log entries.")
                        .build();
                    responseObserver.onNext(summary);
                    responseObserver.onCompleted();
                    System.out.println("Finished processing client stream. Sent summary.");
                }
            };
        }
    }
}

Asiakas: suoratoiston lähettäminen

Asiakaspuolella saatte StreamObserver-olion viestien lähettämistä varten. Kutsutte onNext()-metodia jokaiselle lähetettävälle viestille ja lopuksi onCompleted()-metodia ilmoittaaksenne suoratoiston päättymisestä.

Palvelimen yksittäinen vastaus käsitellään erillisellä antamallanne StreamObserver-oliolla.

import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.stub.StreamObserver;
import client_streaming.LogEntry;
import client_streaming.LogServiceGrpc;
import client_streaming.UploadSummary;

import java.util.concurrent.TimeUnit;

public class LogClient {
    private final LogServiceGrpc.LogServiceStub asyncStub;
    private final ManagedChannel channel;

    public LogClient(String host, int port) {
        channel = ManagedChannelBuilder.forAddress(host, port)
            .usePlaintext() // For demonstration, use plaintext
            .build();
        asyncStub = LogServiceGrpc.newStub(channel);
    }

    public void shutdown() throws InterruptedException {
        channel.shutdown().awaitTermination(5, TimeUnit.SECONDS);
    }

    public void uploadMultipleLogs() throws InterruptedException {
        StreamObserver<UploadSummary> responseObserver = new StreamObserver<UploadSummary>() {
            @Override
            public void onNext(UploadSummary summary) {
                System.out.println("Server Response: " + summary.getStatusMessage() + " (" + summary.getUploadedCount() + " logs)");
            }

            @Override
            public void onError(Throwable t) {
                System.err.println("UploadLogs failed: " + t.getMessage());
            }

            @Override
            public void onCompleted() {
                System.out.println("Server has completed processing.");
            }
        };

        StreamObserver<LogEntry> requestObserver = asyncStub.uploadLogs(responseObserver);

        try {
            // Send multiple log entries
            LogEntry log1 = LogEntry.newBuilder().setMessage("User login attempt").setTimestamp(System.currentTimeMillis()).build();
            LogEntry log2 = LogEntry.newBuilder().setMessage("Database query executed").setTimestamp(System.currentTimeMillis() + 100).build();
            LogEntry log3 = LogEntry.newBuilder().setMessage("API call completed").setTimestamp(System.currentTimeMillis() + 200).build();

            requestObserver.onNext(log1);
            System.out.println("Client sent log 1");
            Thread.sleep(100); // Simulate some delay
            requestObserver.onNext(log2);
            System.out.println("Client sent log 2");
            Thread.sleep(100);
            requestObserver.onNext(log3);
            System.out.println("Client sent log 3");

            // Mark the end of the client stream
            requestObserver.onCompleted();
            System.out.println("Client finished sending logs.");

            // Wait for server response (handled by responseObserver)
            Thread.sleep(1000); // Give time for server to respond
        } catch (RuntimeException e) {
            requestObserver.onError(e);
            throw e;
        }
    }

    public static void main(String[] args) throws Exception {
        LogClient client = new LogClient("localhost", 50051);
        try {
            client.uploadMultipleLogs();
        } finally {
            client.shutdown();
        }
    }
}

Esimerkin suorittaminen

Voitte nähdä asiakaspuolen suoratoiston toiminnassa seuraavasti:

  1. Kääntäkää ensin .proto-tiedosto tarvittavien Java-luokkien luomiseksi.
  2. Suorittakaa LogServer-sovellus. Se alkaa kuunnella pyyntöjä.
  3. Suorittakaa sitten LogClient-sovellus. Se lähettää kolme lokimerkintää ja odottaa palvelimen yhteenvetoa.

Tarkastelkaa sekä asiakkaan että palvelimen konsolitulostetta viestien kulun ymmärtämiseksi.

StreamObserverin tärkeät metodit

StreamObserver-rajapinta on keskeinen suoratoistettujen RPC-kutsujen käsittelyssä. Sekä asiakas että palvelin käyttävät tämän rajapinnan toteutuksia.

  • onNext(T value): Kutsutaan jokaisesta suoratoistossa vastaanotetusta viestistä. Asiakas käyttää sitä viestien lähettämiseen ja palvelin niiden vastaanottamiseen.
  • onError(Throwable t): Kutsutaan, jos RPC-kutsu epäonnistuu tai peruutetaan.
  • onCompleted(): Kutsutaan, kun suoratoisto on päättynyt. Asiakas kutsuu sitä lähetettyään kaikki viestit, ja palvelin kutsuu sitä lähetettyään viimeisen vastauksensa.

Milloin asiakaspuolen suoratoistoa käytetään

Asiakaspuolen suoratoisto sopii tilanteisiin, joissa asiakkaan on lähetettävä palvelimelle suuri määrä tietoa tai sarja toisiinsa liittyviä viestejä ja asiakas tarvitsee vain yhden lopullisen tuloksen.

  • Suurten tiedostojen lataaminen: Tiedoston lähettäminen osissa.
  • Lokien kokoaminen: Asiakas lähettää useita lokimerkintöjä keskitetylle lokipalvelulle.
  • Erätoiminnot: Käsiteltävien kohteiden lähettäminen yhtenä eränä ja yhteenvedon vastaanottaminen.
  • Anturitietojen kerääminen: Mittaustietojen jatkuva lähettäminen anturilaitteesta.

Pikatarkistus

Suunnittelette gRPC-palvelua, jossa asiakkaan on lähetettävä palvelimelle sarja anturimittauksia ja palvelin käsittelee ne ja palauttaa yhden yhteenvetoraportin.

Kertaus: asiakaspuolen suoratoisto

Olette oppineet gRPC:n asiakaspuolen suoratoistosta!

  • Asiakaspuolen suoratoiston avulla asiakas voi lähettää viestisarjan.
  • Palvelin käsittelee viestit ja lähettää takaisin yhden vastauksen.
  • Protobufissa se määritetään käyttämällä stream-avainsanaa pyyntötyypille.
  • Sekä asiakas että palvelin käyttävät StreamObserver-oliota viestien kulun hallintaan (onNext(), onError(), onCompleted()).
  • Se sopii erinomaisesti esimerkiksi suurten tietomäärien lataamiseen tai jatkuvien lokimerkintöjen lähettämiseen.

Seuraavaksi tutustumme palvelinpuolen suoratoistoon!

Aloita maksutta

Opi gRPC ja suorituskykyiset API:t tekoälytuutorin avulla — ilmaiseksi

Kirjoita ja suorita oikeaa koodia selaimessa, saa välitöntä apua tekoälytuutorilta ympäri vuorokauden ja jatka siitä, mihin jäit, verkossa tai sovelluksessa.

Kurssit
12
Oppitunnit
48

Usein kysytyt kysymykset

Onko oppitunti ”Asiakaspohjainen streamaus” ilmainen?

Kyllä – oppitunnin ”Asiakaspohjainen streamaus” koko tekstin voi lukea täällä verkossa ilmaiseksi. Jos haluat harjoitella interaktiivisesti sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla sekä avata koko gRPC ja suorituskykyiset API:t-kurssin, päivitä CoddyKit PROhon. gRPC ja suorituskykyiset API:t-kurssilla on yhteensä 4 oppituntia.

Mitä opin oppitunnilla ”Asiakaspohjainen streamaus”?

Opi toteuttamaan asiakaspään streamaus, jossa asiakas voi lähettää palvelimelle viestisarjan. Harjoittelet gRPC ja suorituskykyiset API:t-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.

Tarvitsenko kokemusta aloittaakseni gRPC ja suorituskykyiset API:t-opiskelun?

Aiempi kokemus ei ole tarpeen. CoddyKitin gRPC ja suorituskykyiset API:t-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 2/4.

Kuinka kauan ”Asiakaspohjainen streamaus”-oppitunnin suorittaminen kestää?

Useimmat CoddyKitin oppitunnit kestävät noin 5–10 minuuttia. Jokainen oppitunti on lyhyt ja interaktiivinen, joten edistyt tasaisesti ja voit jatkaa siitä, mihin jäit – sekä verkossa että sovelluksessa.

Voinko kirjoittaa ja suorittaa koodia tällä gRPC ja suorituskykyiset API:t-oppitunnilla?

Kyllä. Jokainen gRPC ja suorituskykyiset API:t-oppitunti sisältää sisäänrakennetun koodieditorin, joten voit kirjoittaa ja suorittaa oikeaa koodia suoraan selaimessa ja saada välitöntä palautetta tekoälyltä – paikallista asennusta ei tarvita.

Kaikki tämän kurssin oppitunnit

  1. Palvelinpohjainen streamaus
  2. Asiakaspohjainen streamaus
  3. Kaksisuuntainen streamaus
  4. Suoratoiston vuonhallinta ja backpressure
← Takaisin: gRPC ja suorituskykyiset API:t