Asiakaspohjainen streamaus
Opi toteuttamaan asiakaspään streamaus, jossa asiakas voi lähettää palvelimelle viestisarjan.
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:
- Kääntäkää ensin
.proto-tiedosto tarvittavien Java-luokkien luomiseksi. - Suorittakaa LogServer-sovellus. Se alkaa kuunnella pyyntöjä.
- 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!
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
- Palvelinpohjainen streamaus
- Asiakaspohjainen streamaus
- Kaksisuuntainen streamaus
- Suoratoiston vuonhallinta ja backpressure