In diesem Dokument wird beschrieben, wie Sie Daten aus BigQuery in Dataflow lesen.
Übersicht
In den meisten Anwendungsfällen empfiehlt es sich, verwaltete E/A-Vorgänge zum Lesen aus BigQuery zu verwenden. Verwaltete E/A-Vorgänge bieten Funktionen wie automatische Upgrades und eine konsistente Konfigurations-API. Beim Lesen aus BigQuery führt die verwaltete E/A-Funktion direkte Tabellenlesevorgänge aus, die die beste Leseleistung bieten.
Wenn Sie eine erweiterte Leistungsoptimierung benötigen, können Sie den BigQueryIO-Connector verwenden. Der BigQueryIO-Connector unterstützt sowohl direkte Tabellenlesevorgänge als auch das Lesen aus BigQuery-Exportjobs. Außerdem bietet er eine detailliertere Steuerung der Deserialisierung von Tabellendatensätzen. Weitere Informationen finden Sie in diesem Dokument unter
Connector verwendenBigQueryIO.
Spaltenprojektion und -filterung
Um die Datenmenge zu reduzieren, die Ihre Pipeline aus BigQuery liest, können Sie die folgenden Techniken verwenden:
- Bei der Spaltenprojektion wird eine Teilmenge von Spalten angegeben, die aus der Tabelle gelesen werden sollen. Verwenden Sie die Spaltenprojektion, wenn Ihre Tabelle eine große Anzahl von Spalten enthält und Sie nur eine Teilmenge davon lesen müssen.
- Bei der Zeilenfilterung wird ein Prädikat angegeben, das auf die Tabelle angewendet werden soll. Der BigQuery-Lesevorgang gibt nur Zeilen zurück, die dem Filter entsprechen. Dadurch kann die Gesamtmenge der von der Pipeline aufgenommenen Daten reduziert werden.
Im folgenden Beispiel werden die Spalten "user_name" und "age" aus einer Tabelle gelesen
und Zeilen herausgefiltert, die nicht mit dem Prädikat "age > 18" übereinstimmen. In diesem Beispiel wird die verwaltete E/A-Funktion verwendet.
Java
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Dataflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
Aus einem Abfrageergebnis lesen
Im folgenden Beispiel wird die verwaltete E/A-Funktion verwendet, um das Ergebnis einer SQL-Abfrage zu lesen. Dabei wird eine Abfrage für ein öffentliches BigQuery-Dataset ausgeführt. Sie können auch SQL-Abfragen verwenden, um aus einer BigQuery-Ansicht oder einer materialisierten Ansicht zu lesen.
Java
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Dataflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
BigQueryIO-Connector verwenden
Der BigQueryIO-Connector unterstützt die folgenden Serialisierungsmethoden:
- Lesen Sie die Daten als Avro-formatierte Datensätze. Bei dieser Methode stellen Sie eine Funktion bereit, die die Avro-Datensätze in einen benutzerdefinierten Datentyp parst.
- Lesen Sie die Daten als
TableRowObjekte. Diese Methode ist praktisch, da kein benutzerdefinierter Datentyp erforderlich ist. Allerdings ist die Leistung im Allgemeinen geringer als beim Lesen von Avro-formatierten Datensätzen.
Der Connector unterstützt zwei Optionen zum Lesen von Daten:
- Export job. Standardmäßig führt der
BigQueryIOConnector einen BigQuery Exportjob aus, der die Tabellendaten in Cloud Storage schreibt. Der Connector liest dann die Daten aus Cloud Storage. - Direkte Tabellenlesevorgänge. Diese Option ist schneller als Exportjobs, da sie
die BigQuery Storage Read API verwendet und
den Exportschritt überspringt. Wenn Sie direkte Tabellenlesevorgänge verwenden möchten, rufen Sie beim Erstellen der Pipeline
withMethod(Method.DIRECT_READ)auf.
Bei der Auswahl der zu verwendenden Option sollten Sie Folgendes berücksichtigen:
Im Allgemeinen empfehlen wir die Verwendung direkter Tabellenlesevorgänge. Die Storage Read API eignet sich besser für Datenpipelines als Exportjobs, da der Zwischenschritt des Datenexports nicht erforderlich ist.
Wenn Sie direkte Lesevorgänge verwenden, werden Ihnen die Kosten für die Nutzung der Storage Read API in Rechnung gestellt. Weitere Informationen finden Sie auf der Seite BigQuery-Preise unter Preise für die Datenextraktion.
Für Exportjobs fallen keine zusätzlichen Kosten an. Allerdings gelten für Exportjobs Limits. Bei großen Datenmengen, bei denen die Aktualität Priorität hat und die Kosten angepasst werden können, werden direkte Lesevorgänge empfohlen.
Für die Storage Read API gelten Kontingentlimits. Verwenden Sie Google Cloud Messwerte um Ihre Kontingentnutzung zu beobachten.
Wenn Sie Exportjobs verwenden, legen Sie mit der
--tempLocationPipeline-Option einen Cloud Storage-Bucket für die exportierten Dateien fest.Bei Verwendung der Storage Read API werden in den Logs möglicherweise Fehler aufgrund von Lease-Ablauf und Sitzungs-Timeout angezeigt, z. B.:
DEADLINE_EXCEEDEDServer UnresponsiveStatusCode.FAILED_PRECONDITION details = "there was an error operating on 'projects/<projectID>/locations/<location>/sessions/<sessionID>/streams/<streamID>': session`
Diese Fehler können auftreten, wenn ein Vorgang länger als das Timeout dauert, in der Regel in Pipelines, die länger als 6 Stunden ausgeführt werden. Um dieses Problem zu beheben, wechseln Sie zu Dateiexporten.
Der Grad der Parallelität hängt von der Lesemethode ab:
Direkte Lesevorgänge: Der E/A-Connector erzeugt eine dynamische Anzahl von Streams, basierend auf der Größe der Exportanfrage. Diese Streams werden parallel direkt aus BigQuery gelesen.
Exportjobs: BigQuery bestimmt, wie viele Dateien in Cloud Storage geschrieben werden sollen. Die Anzahl der Dateien hängt von der Abfrage und der Datenmenge ab. Der E/A-Connector liest die exportierten Dateien parallel.
Die folgende Tabelle enthält Leistungsmesswerte für verschiedene BigQuery-E/A-Leseoptionen. Die Arbeitslasten wurden auf einem e2-standard2-Worker mit dem Apache Beam SDK 2.49.0 für Java ausgeführt. Der Portable Runner wurde nicht verwendet.
| 100 Mio. Datensätze | 1 KB | 1 Spalte | Durchsatz (Byte) | Durchsatz (Elemente) |
|---|---|---|
| Speicherlesevorgänge | 120 Mbit/s | 88.000 Elemente pro Sekunde |
| Avro-Export | 105 Mbit/s | 78.000 Elemente pro Sekunde |
| JSON-Export | 110 Mbit/s | 81.000 Elemente pro Sekunde |
Diese Messwerte basieren auf einfachen Batch-Pipelines. Sie dienen zum Vergleich der Leistung zwischen E/A-Anschlüssen und sind nicht unbedingt repräsentativ für reale Pipelines. Die Leistung der Dataflow-Pipeline ist komplex und eine Funktion des VM-Typs, der verarbeiteten Daten, der Leistung externer Quellen und Senken sowie des Nutzercodes. Die Messwerte basieren auf der Ausführung des Java SDK und sind nicht repräsentativ für die Leistungsmerkmale anderer Sprach-SDKs. Weitere Informationen finden Sie unter Beam E/A Leistung.
Beispiele
In den folgenden Codebeispielen wird der BigQueryIO-Connector mit direkten Tabellenlesevorgängen verwendet. Wenn Sie stattdessen einen Exportjob verwenden möchten, lassen Sie den Aufruf von withMethod weg.
Avro-formatierte Datensätze lesen
In diesem Beispiel wird gezeigt, wie Sie den BigQueryIO-Connector verwenden, um Avro-formatierte Datensätze zu lesen.
Verwenden Sie die
read(SerializableFunction) Methode, um BigQuery-Daten in Avro-formatierte Datensätze zu lesen. Diese Methode
verwendet eine anwendungsdefinierte Funktion, die
SchemaAndRecord-Objekte parst und einen
benutzerdefinierten Datentyp zurückgibt. Die Ausgabe des Connectors ist eine PCollection Ihres benutzerdefinierten Datentyps.
Der folgende Code liest ein PCollection<MyData> aus einer BigQuery-Tabelle, wobei MyData eine anwendungsdefinierte Klasse ist.
Java
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Dataflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
Die read Methode verwendet eine SerializableFunction<SchemaAndRecord, T> Schnittstelle,
die eine Funktion zum Konvertieren von Avro-Datensätzen in eine benutzerdefinierte Datenklasse definiert. Im vorherigen Codebeispiel implementiert die Methode MyData.apply diese Konvertierungsfunktion. Die Beispielfunktion parst die Felder name und age aus dem Avro-Datensatz und gibt eine MyData-Instanz zurück.
Geben Sie die BigQuery-Tabelle an, die gelesen werden soll, indem Sie die Methode from aufrufen, wie im vorherigen Beispiel gezeigt. Weitere Informationen finden Sie in der Dokumentation zum BigQuery-E/A-Connector unter
Tabellennamen.
TableRow-Objekte lesen
In diesem Beispiel wird gezeigt, wie Sie den BigQueryIO-Connector verwenden, um TableRow-Objekte zu lesen.
Die readTableRows Methode liest
BigQuery-Daten in eine PCollection von
TableRow-Objekten. Jede TableRow ist eine Map von Schlüssel/Wert-Paaren, die eine einzelne Zeile mit Tabellendaten enthält. Geben Sie die BigQuery-Tabelle an, die gelesen werden soll, indem Sie die Methode from aufrufen.
Mit dem folgenden Code wird PCollection<TableRows> aus einer BigQuery-Tabelle gelesen.
Java
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Dataflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
In diesem Beispiel wird auch gezeigt, wie Sie auf die Werte aus dem TableRow-Wörterbuch zugreifen.
Ganzzahlige Werte werden als Strings codiert, um dem exportierten JSON-Format von BigQuery zu entsprechen.