הזרמת שיחות בספריות לקוח של Cloud ל-Java

שיחות סטרימינג מאפשרות דפוסי אינטראקציה מורכבים יותר מאשר בקשות ותשובות פשוטות, ומאפשרות לשלוח או לקבל כמה הודעות בחיבור יחיד.

ספריות לקוח של Cloud Java תומכות בשלושה סוגים של קריאות סטרימינג:

  • סטרימינג מהשרת: השרת שולח לכם זרם של תשובות.
  • הזרמת נתונים מהלקוח: אתם שולחים לשרת זרם של בקשות.
  • סטרימינג דו-כיווני: אתם יכולים לשלוח לשרת זרם של בקשות, והשרת יכול לשלוח לכם בחזרה זרם של תשובות.

הטמעות הסטרימינג מבוססות על הטמעות gRPC-Java לסטרימינג בשרת, בלקוח ובשני הכיוונים.

תמיכה בסטרימינג בכל אמצעי התקשורת

יש תמיכה מלאה בסטרימינג כשמשתמשים ב-gRPC, אבל רק תמיכה חלקית ב-HttpJson. בטבלה הבאה מפורטות הפלטפורמות שבהן אפשר להפעיל סטרימינג.

סוג השידור gRPC HttpJson
סטרימינג מהשרת נתמך נתמך
סטרימינג בצד הלקוח נתמך לא נתמך
סטרימינג דו-כיווני נתמך לא נתמך

שיחות unary (לא סטרימינג) נתמכות גם ב-gRPC וגם ב-HttpJson.

קביעת סוג הסטרימינג

כדי לקבוע את סוג הסטרימינג של השיחה, בודקים את הערך של Callable type שמוחזר:

  • ‫ServerStreamingCallable: סטרימינג מהשרת.
  • ‫ClientStreamingCallable: סטרימינג ללקוח.
  • ‫BidiStreamingCallable: שידור דו-כיווני.

לדוגמה, באמצעות Java-Aiplatform ו-Java-Speech:

// Server Streaming
ServerStreamingCallable<ReadTensorboardBlobDataRequest, ReadTensorboardBlobDataResponse> callable = aiplatformClient.readTensorboardBlobDataCallable();

// Bidirectional Streaming
BidiStreamingCallable<StreamingRecognizeRequest, StreamingRecognizeResponse> callable = speechClient.streamingRecognizeCallable();

ביצוע שיחות בסטרימינג

הדרך שבה מבצעים קריאה לסטרימינג שונה בהתאם לסוג הסטרימינג שבו משתמשים: סטרימינג בצד השרת או סטרימינג דו-כיווני.

סטרימינג מהשרת

אין צורך בהטמעה נוספת לסטרימינג מהשרת. המחלקות ServerStream מאפשרות לבצע איטרציה על זרם התשובות. בדוגמה הבאה מוצג אופן השימוש ב-Java-Maps-Routing כדי להפעיל את Server Streaming API:

try (RoutesClient routesClient = RoutesClient.create()) {
  ServerStreamingCallable<ComputeRouteMatrixRequest, RouteMatrixElement> computeRouteMatrix =
    routesClient.computeRouteMatrixCallable();  
  ServerStream<RouteMatrixElement> stream = computeRouteMatrix.call(
    ComputeRouteMatrixRequest.newBuilder().build());
  for (RouteMatrixElement element : stream) {
    // Do something with response
  }
}

בדוגמה הזו, הלקוח שולח ComputeRouteMatrixRequest יחיד ומקבל זרם של תגובות.

סטרימינג דו-כיווני

כדי להתקשר באמצעות סטרימינג דו-כיווני, צריך לבצע הטמעה נוספת. בדוגמה הבאה מוסבר איך להשתמש ב-Java-Speech כדי להטמיע קריאה לסטרימינג דו-כיווני.

קודם מטמיעים את הממשק ResponseObserver באמצעות הקוד הבא כהנחיה:

class BidiResponseObserver<T> implements ResponseObserver<T> {
  private final List<T> responses = new ArrayList<>();
  private final SettableApiFuture<List<T>> future = SettableApiFuture.create();

  @Override
  public void onStart(StreamController controller) {
    // no-op
  }

  @Override
  public void onResponse(T response) {
    responses.add(response);
  }

  @Override
  public void onError(Throwable t) {
    future.setException(t);
  }

  @Override
  public void onComplete() {
    future.set(responses);
  }

  public SettableApiFuture<List<T>> getFuture() {
    return future;
  }
}

לאחר מכן, פועלים לפי השלבים הבאים:

  1. יוצרים מופע של האובייקט observer:

    BidiResponseObserver<StreamingRecognizeResponse> responseObserver = new BidiResponseObserver<>();
    
  2. מעבירים את הישות שזיהתה את האירוע אל הפונקציה שאפשר להפעיל:

    ClientStream<EchoRequest> clientStream = speechClient.streamingRecognizeCallable().splitCall(responseObserver);
    
  3. שליחת הבקשות לשרת וסגירת הסטרימינג בסיום:

    clientStream.send(StreamingRecognizeRequest.newBuilder().build());
    clientStream.send(StreamingRecognizeRequest.newBuilder().build());
    // ... other requests ...
    clientStream.send(StreamingRecognizeRequest.newBuilder().build());
    clientStream.closeSend();
    
  4. חוזרים על התהליך עם התשובות:

    List<StreamingRecognizeResponse> responses = responseObserver.getFuture().get();
    
    for (StreamingRecognizeResponse response : responses) {
      // Do something with response
    }
    

שגיאות בסטרימינג שלא נתמך

בספריות לקוח שתומכות גם ב-gRPC וגם בהעברות HTTP/JSON, יכול להיות שתגדירו בטעות את ספריית הלקוח להפעלה של קריאת סטרימינג שלא נתמכת. לדוגמה, בהגדרה הבאה מוצג לקוח HttpJson של Java-Speech שמבצע שיחת סטרימינג דו-כיוונית:

// SpeechClient is configured to use HttpJson
try (SpeechClient speechClient = SpeechClient.create(SpeechSettings.newHttpJsonBuilder().build())) {
  // Bidi Callable is not supported in HttpJson
  BidiStreamingCallable<StreamingRecognizeRequest, StreamingRecognizeResponse> callable = speechClient.streamingRecognizeCallable();
  ...
}

השגיאה הזו לא גורמת לשגיאת קומפילציה, אבל היא מופיעה כשגיאת זמן ריצה:

Not implemented: streamingRecognizeCallable(). REST transport is not implemented for this method yet.
Important: The client library MUST be configured with gRPC to use client or bidirectional streaming.