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 create_sink() függvény egy eseménystreamelési szolgáltatásba, például az Apache Kafkába vagy Azure Event Hubs vagy egy deklaratív folyamatból egy Delta-táblába ír. Miután létrehozott egy fogadót a create_sink() függvénnyel, a fogadót hozzáfűző folyamatban vagy frissítési folyamatban használva adatokat írhat a fogadóba. A függvény csak a hozzáfűzési és frissítési create_sink() folyamatokat támogatja. Más folyamattípusok, például create_auto_cdc_flow, nem támogatottak. A Lakeflow-folyamatok más típusú fogadóiról további információt a Lakeflow-folyamatok fogadói című témakörben talál.
A Delta sink támogatja a Unity Catalog külső és felügyelt tábláit, valamint a Hive metaadattár által felügyelt táblákat. A táblaneveknek teljesen specifikáltnak kell lenniük. A Unity Catalog tábláinak például háromszintű azonosítót kell használniuk: <catalog>.<schema>.<table>. A Hive metaadattártábláinak <schema>.<table>kell használniuk.
Megjegyzés:
- A teljes frissítés futtatása nem törli az adatokat a tárolóhelyekről. Az újrafeldolgozott adatok hozzá vannak fűzve a fogadóhoz, és a meglévő adatok nem módosulnak.
- Az API nem támogatja az
sinkelvárásokat.
Szemantika
from pyspark import pipelines as dp
dp.create_sink(name=<sink_name>, format=<format>, options=<options>)
Paraméterek
| Paraméter | Típus | Description |
|---|---|---|
name |
str |
Szükséges. Egy karaktersor, amely azonosítja a csatornát, és a csatorna hivatkozására és kezelésére szolgál. A fogadó neveknek egyedieknek kell lenniük az adatfolyamban, beleértve az adatfolyam részét képező összes forráskódfájlt is. |
format |
str |
Szükséges. A kimeneti formátumot definiáló sztring, kafka vagy delta. |
options |
dict |
A fogadó beállításainak listája, formázva {"key": "value"}, ahol a kulcs és az érték egyaránt sztring. A Kafka és a Delta csatlakozási pontok által támogatott összes Databricks-futtatókörnyezeti beállítás támogatott.
|
Példák
from pyspark import pipelines as dp
# Create a Kafka sink
dp.create_sink(
"my_kafka_sink",
"kafka",
{
"kafka.bootstrap.servers": "host:port",
"topic": "my_topic"
}
)
# Create an external Delta table sink with a file path
dp.create_sink(
"my_delta_sink",
"delta",
{ "path": "/path/to/my/delta/table" }
)
# Create a Delta table sink using a table name
dp.create_sink(
"my_delta_sink",
"delta",
{ "tableName": "my_catalog.my_schema.my_table" }
)