Daten im Batch mit der Storage Write API (gRPC) laden

In diesem Dokument wird beschrieben, wie Sie die BigQuery Storage Write API (gRPC) verwenden, um Daten im Batch in BigQuery zu laden.

In Szenarien mit Batchladevorgängen schreibt eine Anwendung Daten und übergibt sie per Commit in einer einzigen unteilbaren Transaktion. Wenn Sie die Storage Write API (gRPC) für das Laden von Daten im Batch verwenden, erstellen Sie einen oder mehrere Streams vom Typ „Ausstehend“. Der Typ „Ausstehend“ unterstützt Transaktionen auf Streamebene. Datensätze werden im Status „Ausstehend“ zwischengespeichert, bis Sie den Stream per Commit übergeben.

Prüfen Sie bei Batcharbeitslasten auch die Verwendung der Storage Write API (gRPC) über den Apache Spark SQL Connector für BigQuery mithilfe von Managed Service for Apache Spark, anstatt benutzerdefinierten Storage Write API (gRPC)-Code zu schreiben.

Die Storage Write API (gRPC) eignet sich gut für eine Datenpipeline-Architektur. Ein Hauptprozess erzeugt eine Reihe von Streams. Jedem Stream wird ein Worker-Thread oder ein separater Prozess zugewiesen, um einen Teil der Batch-Daten zu schreiben. Jeder Worker erstellt eine Verbindung zu seinem Stream, schreibt Daten und finalisiert seinen Stream, wenn er abgeschlossen ist. Nachdem alle Worker den erfolgreichen Abschluss zum Hauptprozess signalisiert haben, übergibt der Hauptprozess die Daten per Commit. Wenn ein Worker fehlschlägt, wird der zugewiesene Teil der Daten nicht in den Endergebnissen angezeigt und der gesamte Worker kann wiederholt werden. In einer komplexeren Pipeline prüfen Worker ihren Fortschritt, indem sie den letzten an den Hauptprozess geschriebenen Offset melden. Dieser Ansatz kann zu einer robusten Pipeline führen, die ausfallsicher ist.

Daten im Batch unter Verwendung des Typs „Ausstehend“ laden

Die Anwendung geht so vor, um den Typ „Ausstehend“ zu verwenden:

  1. Rufen Sie CreateWriteStream auf, um einen oder mehrere Streams vom Typ „Ausstehend“ zu erstellen.
  2. Rufen Sie für jeden Stream AppendRows in einer Schleife auf, um Datensätze in Batches zu schreiben.
  3. Rufen Sie für jeden Stream FinalizeWriteStream auf. Nach dem Aufrufen dieser Methode können Sie keine weiteren Zeilen in den Stream schreiben. Wenn Sie AppendRows nach dem Aufruf von FinalizeWriteStream aufrufen, wird ein StorageError mit StorageErrorCode.STREAM_FINALIZED im Fehler google.rpc.Status zurückgegeben. Weitere Informationen zum Fehlermodell google.rpc.Status finden Sie unter Fehler.
  4. Rufen Sie BatchCommitWriteStreams auf, um die Streams per Commit zu übergeben. Nach dem Aufrufen dieser Methode stehen die Daten zum Lesen zur Verfügung. Wenn beim Commit eines der Streams ein Fehler auftritt, wird der Fehler im Feld stream_errors der BatchCommitWriteStreamsResponse zurückgegeben.

Das Commit ist ein unteilbarer Vorgang und Sie können mehrere Streams gleichzeitig per Commit übergeben. Für einen Stream kann nur einmal ein Commit durchgeführt werden. Wenn der Commit-Vorgang fehlschlägt, kann der Vorgang sicher wiederholt werden. Bis zum Commit eines Streams stehen die Daten aus und sind für Lesevorgänge nicht sichtbar.

Nachdem der Stream abgeschlossen wurde und bevor er übergeben wird, können die Daten bis zu 4 Stunden im Puffer verbleiben. Ausstehende Streams müssen innerhalb von 24 Stunden per Commit bestätigt werden. Für die Gesamtgröße des Zwischenspeichers für ausstehende Streams gilt ein Kontingentlimit.

Der folgende Code zeigt, wie Daten vom Typ „Ausstehend“ geschrieben werden.

C#

Informationen zum Installieren und Verwenden der Clientbibliothek für BigQuery finden Sie unter BigQuery-Clientbibliotheken. Weitere Informationen finden Sie in der Referenzdokumentation zur BigQuery C# API.

Richten Sie zur Authentifizierung bei BigQuery die Standardanmeldedaten für Anwendungen ein. Weitere Informationen finden Sie unter Authentifizierung für Clientbibliotheken einrichten.


using Google.Api.Gax.Grpc;
using Google.Cloud.BigQuery.Storage.V1;
using Google.Protobuf;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using static Google.Cloud.BigQuery.Storage.V1.AppendRowsRequest.Types;

public class AppendRowsPendingSample
{
    /// <summary>
    /// This code sample demonstrates how to write records in pending mode.
    /// Create a write stream, write some sample data, and commit the stream to append the rows.
    /// The CustomerRecord proto used in the sample can be seen in Resources folder and generated C# is placed in Data folder in
    /// https://github.com/GoogleCloudPlatform/dotnet-docs-samples/tree/main/bigquery-storage/api/BigQueryStorage.Samples
    /// </summary>
    public async Task AppendRowsPendingAsync(string projectId, string datasetId, string tableId)
    {
        BigQueryWriteClient bigQueryWriteClient = await BigQueryWriteClient.CreateAsync();
        // Initialize a write stream for the specified table.
        // When creating the stream, choose the type. Use the Pending type to wait
        // until the stream is committed before it is visible. See:
        // https://cloud.google.com/bigquery/docs/reference/storage/rpc/google.cloud.bigquery.storage.v1#google.cloud.bigquery.storage.v1.WriteStream.Type
        WriteStream stream = new WriteStream { Type = WriteStream.Types.Type.Pending };
        TableName tableName = TableName.FromProjectDatasetTable(projectId, datasetId, tableId);

        stream = await bigQueryWriteClient.CreateWriteStreamAsync(tableName, stream);

        // Initialize streaming call, retrieving the stream object
        BigQueryWriteClient.AppendRowsStream rowAppender = bigQueryWriteClient.AppendRows();

        // Sending requests and retrieving responses can be arbitrarily interleaved.
        // Exact sequence will depend on client/server behavior.
        // Create task to do something with responses from server.
        Task appendResultsHandlerTask = Task.Run(async () =>
        {
            AsyncResponseStream<AppendRowsResponse>