Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
A Structured Streaming használatával adatokat írhat a Lakebase-be vagy egy külső PostgreSQL-adatbázisba beépített kötegeléssel, automatikus újrapróbálkozásokkal és munkaterület által kezelt hitelesítéssel.
Mikor használja a Lakebase sinket?
Használd a Lakebase sinket alacsony késleltetésű streaming írásokhoz Lakebase-be vagy egy külső PostgreSQL adatbázisba. Ez a fogadó nem követeli meg, hogy egyéni foreach függvényeket implementáljon a kötegelés, a kapcsolatkezelés és a hibakezelés kezeléséhez.
Gyakori használati esetek a következők:
- Valós időben frissítheti az alkalmazás-adatbázisokat az operatív irányítópultok vagy az ügyféloldali funkciók számára.
- Szinkronizálja a folyamatosan változó adatokat, például az összesített vagy szűrt streamelési eredményeket egy tranzakciós adatbázisba.
- Írja be egy strukturált streamelési lekérdezés kimenetét egy Lakebase-táblába, amely másodperc alatti késéssel rendelkezik valós idejű módban.
A Lakebase-ből a Lakehouse Delta Lake-tábláiba történő adatszinkronizáláshoz, a fordított iránnyal kapcsolatban lásd: Lakebase Change Data Feed.
Követelmények
-
Databricks Runtime 18 LTS és annál magasabb.
- Külső PostgreSQL kapcsolatokhoz a Databricks Runtime 19 vagy annál magasabb verziókat kell használnod, és csatlakoznod a UC Compute előnézeten található Custom JDBC-hez .
- Az intervallum adattípusok megkövetelik, hogy a Databricks Runtime 19 vagy annál magasabb verziókat használj.
- Klasszikus számítástechnika dedikált vagy szabványos hozzáférési módokkal, vagy szerver nélküli számítás notebookokhoz vagy feladatokhoz. Szerver nélküli számítás esetén használd
Trigger.AvailableNow(). Lásd : Streamelés kiszolgáló nélküli számításon. - Egy Lakebase adatbázis, vagy egy Unity Catalog kapcsolat egy külső PostgreSQL adatbázishoz.
Azonosító követelmények
Minden célpontnál a Databricks azt javasolja, hogy olyan séma, táblázat, oszlop és elsődleges kulcs oszlopnevek legyenek, amelyek betűvel vagy aláhúzással kezdődnek, és csak betűket, számokat és aláhúzásokat tartalmaznak. A mosó ezeket a követelményeket érvényesíti, amikor automatikusan létrehoz egy Lakebase táblát. Ha olyan azonosítókat szeretnél használni, amelyek nem felelnek meg ezeknek a követelményeknek, a lekérdezés elindítása előtt hozd létre a céltáblát.
Kapcsolódás adatbázishoz
A Lakebase-fogadó a következő csatlakozási módszereket támogatja:
A Unity Katalógusban regisztrált Lakebase-táblák
A Unity Catalogban regisztrált Lakebase-táblák esetében az összekötő automatikusan kezeli a hitelesítő adatokat, és a lekérdezést futtató felhasználó vagy szolgáltatásnév identitását használja. Ha a tábla nem létezik, az összekötő létrehozza a táblát.
Ha lakebase-adatbázist szeretne regisztrálni a Unity Catalogban, olvassa el a Lakebase-adatbázis regisztrálása a Unity Catalogban című témakört.
Egy Lakebase táblához írni egy .toTable() teljesen minősített táblanévű metódusot használjuk: catalog.schema.table
Python
(df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
)
Scala
df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-columns>") // Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
Cserélje ki a következő kitöltendő elemeket.
-
<catalog>.<schema>.<table>: A céltábla teljes neve. Ezcataloga Lakebase-adatbázis regisztrálásakor létrehozott Unity Catalog-katalógus. Lásd: Lakebase-adatbázis regisztrálása a Unity Katalógusban. Ha a tábla nem létezik, az összekötő létrehozza. -
<primary-key-columns>:Szabadon választható. Egy vesszővel elválasztott lista a céltáblás elsődleges kulcsának összes oszlopáról, példáulidvagyuser_id,event_type. Lásd : Upsert viselkedés. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Egy Unity-katalógus kötetútvonala, ahol a lekérdezés tárolja az ellenőrzőpontot. Felhőalapú objektumtárolási URI-t is használhat. A helynek olyan tárolónak kell lennie, amelybe írhat, nem pedig helyi lemezre, és egyedinek kell lennie az egyes streamelési lekérdezések esetében. Ez független a céltáblától. Lásd: Strukturált streamelési ellenőrzőpontok.
Opcionális konfigurációkért, mint batchsize például és batchinterval, lásd a PostgreSQL sink opciókat.
A Unity Catalogban nem regisztrált Lakebase-táblák
A Unity Catalogban nem regisztrált Lakebase-táblák esetében az összekötő automatikusan kezeli a hitelesítő adatokat, és a lekérdezést futtató felhasználó vagy szolgáltatásnév identitását használja. Ha a tábla nem létezik, az összekötő létrehozza a táblát.
A Lakebase-táblába való íráshoz használja a(z) endpoint és dbtable opciókat:
Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") # Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)
Scala
df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") // Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") // Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
Cserélje ki a következő kitöltendő elemeket.
-
<project-id>.<branch-id>.<endpoint-id>: A Lakebase-végpont. Keresse meg mindhárom értéket az erőforrásnévben a Számítások lap Get ID menüjében, amelynek formátumaprojects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Lásd : Számítási azonosítók. -
<database>:Szabadon választható. A cél PostgreSQL adatbázis neve. Alapértelmezett érték:databricks_postgres. Lásd: Adatbázisok kezelése. -
<schema>.<table>: A(z)schema.tableformátumú céltábla. Ha kihagyja a sémát, a fogadó apublicsémát használja. Az automatikus táblázatkészítéshez olyan azonosítókat használjunk, amelyek betűvel vagy aláhúzással kezdődnek, és csak betűket, számokat és aláhúzásokat tartalmaznak. -
<primary-key-columns>:Szabadon választható. Egy vesszővel elválasztott lista a céltáblás elsődleges kulcsának összes oszlopáról, példáulidvagyuser_id,event_type. Lásd : Upsert viselkedés. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Egy Unity-katalógus kötetútvonala, ahol a lekérdezés tárolja az ellenőrzőpontot. Felhőalapú objektumtárolási URI-t is használhat. A helynek olyan tárolónak kell lennie, amelybe írhat, nem pedig helyi lemezre, és egyedinek kell lennie az egyes streamelési lekérdezések esetében. Ez független a céltáblától. Lásd: Strukturált streamelési ellenőrzőpontok.
Opcionális konfigurációkért, mint batchsize például és batchinterval, lásd a PostgreSQL sink opciókat.
Külső PostgreSQL Unity Catalog-hitelesítő adatokkal
Important
Ez a funkció nyilvános előzetes verzióban van. A munkaterület adminisztrátorai az Custom JDBC-hez való hozzáférést az UC Compute-on az Előnézetek oldalán irányíthatják. Lásd: Az Azure Databricks előzetes verziójának kezelése.
Használj Unity Catalog kapcsolatot egy külső PostgreSQL adatbázishoz hitelesítéshez anélkül, hogy a kódodban tárolnád a hitelesítő adatokat. A céltáblának már léteznie kell.
Hozzon létre egy típusú POSTGRESQLkapcsolatot, lásd: Kapcsolat létrehozása. A lekérdezést futtató felhasználónak vagy szolgáltatásnévnek rendelkeznie kell USE CONNECTION a kapcsolaton.
A PostgreSQL-táblába való íráshoz használja a database, dbtable és databricks.connection opciókat:
Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("databricks.connection", "<connection-name>")
.option("database", "<database>")
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)
Scala
df.writeStream
.format("postgresql")
.outputMode("update")
.option("databricks.connection", "<connection-name>")
.option("database", "<database>")
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") // Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
Cserélje ki a következő kitöltendő elemeket.
-
<connection-name>: Az Unity Catalog kapcsolatának neve. -
<database>: A cél PostgreSQL adatbázis neve. -
<schema>.<table>: A meglévő céltáblaschema.tableformátumban. Ha kihagyja a sémát, a fogadó apublicsémát használja. -
<primary-key-columns>:Szabadon választható. Egy vesszővel elválasztott lista a céltáblás elsődleges kulcsának összes oszlopáról, példáulidvagyuser_id,event_type. Lásd : Upsert viselkedés. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Egy Unity-katalógus kötetútvonala, ahol a lekérdezés tárolja az ellenőrzőpontot. Felhőalapú objektumtárolási URI-t is használhat. A helynek olyan tárolónak kell lennie, amelybe írhat, nem pedig helyi lemezre, és egyedinek kell lennie az egyes streamelési lekérdezések esetében. Ez független a céltáblától. Lásd: Strukturált streamelési ellenőrzőpontok.
A PostgreSQL kapcsolatok mindig TLS-t használnak. A tanúsítvány ellenőrzése a Unity katalógus kapcsolat beállításait követi, amelyeket a kapcsolat létrehozásakor választasz:
-
Megbízható szerver tanúsítvány: Ha kiválasztva, a kapcsolat használ
sslmode=require, amely titkosítja a kapcsolatot anélkül, hogy ellenőrizné a szerver tanúsítványt. -
A felhasználó által biztosított szervertanúsítvány: Biztosíts egy PEM-kódolt szervertanúsítványt, amit akkor használhatsz
sslmode=verify-full, ha a Trust szerver tanúsítvány nincs kiválasztva. Ha nem adsz tanúsítványt, a kapcsolat a JVM alapértelmezett bizalmi tárolóját használjasslmode=verify-full.
Konfigurációs beállítások
A nyelő hibát jelez az ismeretlen beállítások esetén, JDBC_STREAMING_SINK_INVALID_OPTIONS.
A sink konfigurációs opciók, beleértve a közös opciókat és az egyes csatlakozási módszerek lehetőségeit, lásd a PostgreSQL sink opciókat.
Adattípus-leképezések
Az elnyelő ellenőrzi, hogy minden DataFrame oszlop kompatibilis-e a megfelelő céloszloptal, mielőtt egy meglévő Lakebase vagy külső PostgreSQL táblába írna.
Az alábbi táblázat tartalmazza a Databricks Runtime 18 LTS és annál magasabb rendszerekben támogatott típusokat:
| Spark-típus | Automatikusan létrehozott Lakebase-táblatípus | Kompatibilis típusok meglévő PostgreSQL táblákban |
|---|---|---|
ByteType, ShortType |
smallint |
smallint |
IntegerType |
integer |
integer |
LongType |
bigint |
bigint |
FloatType |
real |
real |
DoubleType |
double precision |
double precision |
DecimalType |
numeric |
numeric |
StringType |
text |
varchar, text |
VarcharType(n) |
varchar(n) |
varchar, text |
CharType(n) |
char(n) |
char |
BinaryType |
bytea |
bytea |
BooleanType |
boolean |
boolean |
TimestampType |
timestamptz |
timestamptz |
TimestampNTZType |
timestamp |
timestamp |
DateType |
date |
date |
ArrayType, MapType, StructType, VariantType, NullType |
jsonb |
json, jsonb |
Az alábbi táblázat tartalmazza a Databricks Runtime 19 és annál magasabb rendszerekben támogatott típusokat:
| Spark-típus | Automatikusan létrehozott Lakebase-táblatípus | Kompatibilis típusok meglévő PostgreSQL táblákban |
|---|---|---|
DayTimeIntervalType, YearMonthIntervalType |
interval |
interval |
Felerősítő viselkedés
Az upsertkey opció azonosítja a céltáblák elsődleges kulcsoszlopait. Egy meglévő tábla esetében az oszlopoknak upsertkey pontosan egyezniük kell a tábla elsődleges kulcsával. Ha kihagyod az opciót, akkor a sink az elsődleges kulcsot a táblából olvassa ki. Egy olyan Lakebase-tábla esetében, amelyet a sink hoz létre, a upsertkey határozza meg az elsődleges kulcsot. Ha kihagyja ezt a beállítást, a sink elsődleges kulcs nélkül hozza létre a táblát.
Ha a céltáblának van elsődleges kulcsa, a sink a PostgreSQL INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... szintaxisával upsert műveletet végez. Ha a céltáblában nincs elsődleges kulcs, a nyelő beszúrásokat végez. A lekérdezés kimeneti módja nem befolyásolja ezt a viselkedést.
Minden elsődleges kulcsoszlopnak jelen kell lennie a DataFrame-ben, és hasonló típusokat kell használni, például numerikus vagy string típusokat.
Teljesítményhangolás
Kötegelés és visszanyomás
A kiürítés akkor aktiválódik, ha valamelyik feltétel teljesül:
- A puffer mérete eléri a
batchsizesort, amelynek alapértelmezett értéke1000. - A puffer kora meghaladja a(z)
batchintervalértéket, amelynek alapértelmezett értéke100 milliseconds.
Ha az adatbázis nem képes lépést tartani a bejövő adatok sebességével, a nyelő backpressure-t továbbít visszafelé a forrás felé.
Késleltetésre és áteresztőképességre vonatkozó útmutató:
- Valós idejű móddal működő, alacsony késleltetésű munkaterhelések esetén csökkentse a(z)
batchintervalértékét, hogy garantálja a kiírás előtti rövidebb maximális időt. A valós idejű mód fogalmakért és valós idejű mód példáit kód példákért lásd. - A nagy átviteli sebességű számítási feladatok esetében növelje
batchsizeaz egyes tranzakciók terhelésének csökkentéséhez.
Kapcsolat viselkedése
A nyelő kapcsolatpoolozást használ a végrehajtó folyamatokon. Alapértelmezés szerint minden tevékenység egy adatbázis-kapcsolatot használ.
A Databricks azt javasolja, hogy minden kapcsolathoz használja a tevékenység alapértelmezett értékét 1 . Ha növeli az egyes kapcsolatokhoz tartozó feladatok számát, előfordulhat, hogy kapcsolati versengéseket okoz, és növeli a késéseket a nagy átviteli sebességű kapcsolatok esetében.
A tevékenységek kapcsolatokhoz viszonyított arányának konfigurálásához állítsa be a Spark-konfigurációt spark.databricks.sql.streaming.jdbc.tasksPerConnection . Ha a céladatbázisnak alacsony a kapcsolati korlátja, csökkentse az shuffle partíciók számát, vagy növelje spark.databricks.sql.streaming.jdbc.tasksPerConnection.
A fogadó automatikusan újrapróbálkozza az átmeneti JDBC-hibákat, beleértve a csatlakozási hibákat, a holtpontokat és a sebességkorlátozást. Ha a nyelőnél kimerül az összes újrapróbálkozási lehetőség, a lekérdezés meghiúsul.
Támogatott eseményindítók és kimeneti módok
Kiváltó okok
Ez a táblázat a strukturált streaminges triggertípusok támogatását mutatja klasszikus és szerver nélküli számítási rendszeren:
| Trigger | Klasszikus számítás | Kiszolgáló nélküli számítási kapacitás (notebookok és feladatok) |
|---|---|---|
RealTime |
Yes | No |
ProcessingTime |
Yes | No |
AvailableNow |
Yes | Yes |
Once |
Yes. Deprecated. Használja a AvailableNow. |
Yes. Deprecated. Használja a AvailableNow. |
Kimeneti módok
Ez a táblázat a strukturált streamelési kimeneti módok támogatását mutatja be:
| Kimeneti mód | Supported |
|---|---|
update |
Yes |
append |
Yes. A viselkedés megegyezik a update. A lekérdezés frissíti a meglévő sort, vagy beszúr egy újat, ha a céltábla elsődleges kulccsal rendelkezik; ellenkező esetben csak beszúr. Lásd : Upsert viselkedés. |
complete |
No |
korlátozások
- Egy külső PostgreSQL adatbázis esetén, amely Unity katalóguskapcsolaton keresztül csatlakozik, a céltáblának már léteznie kell. A sink csak a Lakebase-ben hozza létre automatikusan a hiányzó táblákat.
- A tóvízvezetékek nem támogatottak.