Csatlakozás a Lakebase-hez

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. Ez catalog a 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ául id vagy user_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.table formátumú céltábla. Ha kihagyja a sémát, a fogadó a public sé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ául id vagy user_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ábla schema.table formátumban. Ha kihagyja a sémát, a fogadó a public sé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ául id vagy user_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álja sslmode=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 batchsize sort, amelynek alapértelmezett értéke 1000.
  • A puffer kora meghaladja a(z) batchinterval értéket, amelynek alapértelmezett értéke 100 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 batchsize az 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.