Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Important
Dieses Feature befindet sich in der Betaversion.
Ein REPLACE USING Flow hält eine Zieltabelle synchron mit einer Streaming-Quelle: Er ersetzt alle Zeilen, die mit den angegebenen Schlüsselspalten übereinstimmen, und lässt alle anderen Daten unverändert.
Eine Spalte ordnet SEQUENCE BY die Aktualisierungen so, dass das Ergebnis auch dann korrekt ist, wenn die Updates nicht in der richtigen Reihenfolge eintreffen. Für jeden Schlüssel gewinnt die höchste Sequenz, und eine niedrigere Zeile überschreibt niemals eine höhere, die bereits im Ziel ist. Zeilen, die denselben Schlüssel und dieselbe Sequenz teilen, werden angehängt statt ersetzt.
Wie REPLACE USING funktioniert
Betrachten wir eine Ereignistabelle, die Klick- und Konvertierungsereignisse für zwei Regionen speichert, sequenziert durch seq:
| region_id | Gerätetyp | Event_type | Seq |
|---|---|---|---|
| 1 | iOS | klicken | 1 |
| 1 | Android | Umwandlung | 1 |
| 2 | iOS | klicken | 1 |
| 2 | Desktop | klicken | 1 |
Ein REPLACE USING (region_id) SEQUENCE BY seq Flow erhält diese Updates für die Regionen 1 und 3. Region 2 hat keine Updates:
| region_id | Gerätetyp | Event_type | Seq |
|---|---|---|---|
| 1 | iOS | klicken | 2 |
| 1 | Android | Umwandlung | 2 |
| 1 | Desktop | klicken | 2 |
| 3 | iOS | klicken | 1 |
| 3 | Desktop | klicken | 2 |
Das Ziel wird:
| region_id | Gerätetyp | Event_type | Seq | Ergebnis |
|---|---|---|---|---|
| 1 | iOS | klicken | 2 | Ersetzt, weil seq 2 größer ist als seq 1 |
| 1 | Android | Umwandlung | 2 | Ersetzt, weil seq 2 größer ist als seq 1 |
| 1 | Desktop | klicken | 2 | Ersetzt, weil seq 2 größer ist als seq 1 |
| 2 | iOS | klicken | 1 | Unberührt, weil der Schlüssel in diesem Update nicht vorhanden ist |
| 2 | Desktop | klicken | 1 | Unberührt, weil der Schlüssel in diesem Update nicht vorhanden ist |
| 3 | Desktop | klicken | 2 | Hinzugefügt. Die Seq-1-Zeile für Region 3 wird nicht hinzugefügt, da nur die höchste Folge für eine Taste angewendet wird. |
Anforderungen
ERSETZEN MIT Flows erfüllen folgende Anforderungen:
- ERSETZEN MIT Flows, die auf Databricks Runtime 18.2 und höher laufen, auf klassischer oder serverloser Rechenleistung. Databricks empfiehlt Unity Catalog.
- Die Quelle muss eine Streamingquelle sein. REPLACE USING lehnt eine nicht-streamende Quelle ab.
- Sie müssen mindestens eine Schlüsselspalte und genau eine
SEQUENCE BYSpalte angeben.
Wann man ERSETZEN MIT Flows verwenden sollte
Lakeflow-Pipelines bieten drei Ströme, die bestehende Zeilen überschreiben. Wählen Sie basierend darauf, wie Ihre Quelle aussieht und wie sie die zu ersetzenden Zeilen identifiziert:
- Verwenden Sie REPLACE USING wenn Ihre Quelle eine Reihe von teilweisen Snapshots ist, die per Spalte verschlüsselt sind. REPLACE USING überschreibt nur die Daten, die in den eingehenden Daten übereinstimmen, sodass alle anderen Daten unberührt bleiben. Es braucht keinen Primärschlüssel.
- Verwenden Sie AUTO CDC, wenn Ihre Quelle ein Change Data Capture (CDC)-Feed mit expliziten Einfügungs-, Aktualisierungs- und Löschoperationen ist, oder wenn Sie eine Sage Changing Dimension (SCD) Typ-2-Geschichte benötigen. Die AUTO CDC erfordert außerdem einen echten Primärschlüssel. Siehe Die AUTO CDC-APIs: Vereinfachen der Änderungsdatenerfassung mit Pipelines.
- Verwende VERVANG, WHERE wenn deine Quelle ein Snapshot ist und du einen Bereich der von einem Prädikat ausgewählten Zieltabelle neu berechnen und überschreiben möchtest, zum Beispiel in den letzten 7 Tagen, als Batch-Operation. Es braucht keinen Primärschlüssel. Siehe Batchverarbeitung mit REPLACE-FlüssenWHERE.
Erstelle einen ERSETZEN-VERWENDUNGS-Fluss
Definiere ERSETZEN USING Flows entweder in SQL oder Python.
SQL
Verwenden Sie die Klausel FLOW REPLACE USING inline mit CREATE STREAMING TABLE:
CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);
Alternativ können Sie die Longformsyntax CREATE FLOW verwenden:
CREATE STREAMING TABLE payments_current;
CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);
Note
BY NAME ist in SQL erforderlich. Sie gleicht Spalten anhand des Namens und nicht der Position ab.
Python
Deklarieren Sie die Tabelle und den Fluss zusammen mit @dp.table:
from pyspark import pipelines as dp
@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
return spark.readStream.table("samples.wanderbricks.payments")
Alternativ zielt man auf eine bestehende Streaming-Tabelle mit @dp.replace_flow:
from pyspark import pipelines as dp
dp.create_streaming_table("payments_current")
@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
return spark.readStream.table("samples.wanderbricks.payments")
replace_using ist eine Liste von Schlüsselspalten.
sequence_by ist ein Spaltenname oder ein Column Ausdruck und ist erforderlich, wann immer replace_using gesetzt ist.
Sequenzierung und Out-of-Order-Daten
Die Spalte SEQUENCE BY macht das Ergebnis unabhängig von der Reihenfolge, in der Updates eintreffen. Eine Zeile wird nur dann auf eine Taste angewendet, wenn deren Folge größer ist als die bereits für diesen Schlüssel gespeicherte Sequenz, sodass eine späte oder wiederholte Zeile, die älter als der aktuelle Wert ist, ignoriert wird. Schlüssel, die bei einem Update nicht vorhanden sind, bleiben unberührt.
Befolgen Sie diese Praktiken, damit sich der Ersatz vorhersehbar verhält:
| Praxis | Grund |
|---|---|
| Verwenden Sie eine Sequenz, die sich streng pro Schlüsselversion erhöht, wie zum Beispiel einen Zeitstempel, eine Versionsnummer oder einen Log-Offset. | Zwei Zeilen mit derselben Taste und derselben Reihenfolge werden beide beibehalten, was zu doppelten Zeilen für diesen Schlüssel führt. |
| Verwenden Sie eine nicht-nulle Sequenz. | Eine Nullfolge kann zu undefiniertem Verhalten führen. |
Erwartungshaltung
ERSETZEN MIT Flüssen und unterstützen Erwartungen.
warn
fail und verhält sich wie bei anderen Flows: warn Es verstößt weiterhin gegen Zeilen und zeichnet den Verstoß auf und fail stoppt das Update. Weitere Informationen finden Sie unter Verwalten der Datenqualität mit Pipelineerwartungen.
Eine Erwartung behandelt einen verletzenden Streit so, als hätte drop die Quelle ihn nie produziert. Die weggelassene Zeile ersetzt, löscht oder ändert keine passenden Schlüssel in der Zieltabelle:
- Das Droppen erfolgt vor der Deduplizierung, sodass der Flow die aktuellste gültige Version des Schlüssels behält.
- Wenn jede eingehende Zeile für einen Schlüssel weggelassen wird, bleiben die bestehenden Zeilen des Schlüssels unberührt.
- Da eine weggelassene Zeile keinen Sequenzuntergrund setzt, landet ein späteres, gültiges Update trotzdem, selbst wenn die Sequenz niedriger ist als die der gestrichenen Zeile.
Einschränkungen
ERSETZEN MIT Flüssen hat folgende Einschränkungen:
- REPLACE USING unterstützt einen einzelnen Fluss pro Zieltabelle. Die Kombination von ERSETZEN USING mit einem anderen Flusstyp auf demselben Ziel wird nicht unterstützt.
- Die Zieltabelle muss innerhalb der Pipeline erstellt werden.
- Die Quelle muss eine Streamingquelle sein.
- Sie müssen mindestens eine Schlüsselspalte und eine Spalte
SEQUENCE BYangeben. Schlüsselspalten können nicht wiederholt werden, und der Typ jeder Schlüsselspalte muss sortierbar sein. Atomare Typen wie ganze Zahlen, Zeichenketten und Daten können Schlüssel sein, währendMAPundVARIANTnicht.
Beispiele
Die folgenden Beispiele lesen sich aus samples.wanderbricks.booking_updates, einer Beispieltabelle mit Buchungszustandsänderungen, die in jedem von Unity Catalog unterstützten Arbeitsbereich verfügbar ist. Jede Buchung erscheint einmal pro Änderung und booking_id wiederholt sich mit einem neuen booking_update_id. Siehe Wanderbricks-Datensatz.
Beispiel 1: Führen Sie für jede Taste den neuesten Datensatz
Führe nur den aktuellen Stand jeder Buchung auf. Der Fluss tastet auf booking_id und folgt durch , booking_update_idsodass die aktuellste Aktualisierung für eine Buchung die früheren ersetzt. Nutze stattdessen AUTO CDC, wenn dein Quellcode ein Change-Feed mit expliziten Einfügen-, Update- und Löschoperationen ist.
SQL
CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);
Python
from pyspark import pipelines as dp
@dp.table(
name="bookings_current",
replace_using=["booking_id"],
sequence_by="booking_update_id"
)
def bookings_current():
return spark.readStream.table("samples.wanderbricks.booking_updates")
Dieses Beispiel folgt anstelle des booking_update_id Zeitstempels updated_at , da mehrere Aktualisierungen derselben Buchung einen Zeitstempel teilen können. Reihen, die an der Sequenz zusammenhängen, werden angehängt statt ersetzt, was mehr als eine Reihe für diese Buchungen übrig lassen würde.
Beispiel 2: Schlüssel auf mehr als einer Spalte
Wenn ein Datensatz durch eine Kombination von Spalten identifiziert wird, listen Sie sie alle in REPLACE USINGauf. Hier wird jede Buchung durch (property_id, booking_id)identifiziert, sodass der Fluss den aktuellen Zustand jeder Buchung pro Immobilie beibehält. Wenn eine Schlüsselspalte null sein kann, stimmt REPLACE USING null mit null überein, anstatt die Zeile zu überspringen.
SQL
CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);
Python
from pyspark import pipelines as dp
@dp.table(
name="bookings_by_property",
replace_using=["property_id", "booking_id"],
sequence_by="booking_update_id"
)
def bookings_by_property():
return spark.readStream.table("samples.wanderbricks.booking_updates")
Beispiel 3: Ungültige Datensätze mit einer Erwartung fallen lassen
Füge die Erwartung hinzu, dass schlechte Reihen vom Ziel ferngehalten werden. Eine weggelassene Zeile wird so behandelt, als hätte die Quelle sie nie produziert: Sie ersetzt oder löscht nicht den entsprechenden Schlüssel, und der Fluss fällt auf die zuletzt gültige Zeile für diesen Schlüssel zurück. Dieser Flow lässt Updates fallen, die kein positives total_amountErgebnis haben.
from pyspark import pipelines as dp
@dp.table(
name="bookings_validated",
replace_using=["booking_id"],
sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
return spark.readStream.table("samples.wanderbricks.booking_updates")