Doelen in Lakeflow-pipelines

Pijplijnstromen schrijven standaard resultaten naar Delta-tabellen die worden beheerd door Unity Catalog, meestal streamingtabellen of gerealiseerde weergaven. Sinks zijn een alternatief uitvoerdoel waarmee u getransformeerde gegevens kunt schrijven naar bestemmingen buiten door Databricks beheerde opslag, zoals streamingservices voor gebeurtenissen of aangepaste gegevensarchieven.

Sinks worden gebruikt met toevoegstromen. Je definieert een sink met behulp van een van de sink-API's en verwijst er vervolgens naar als de target in je append_flow-definitie.

Wanneer moet u sinks gebruiken

Databricks raadt het gebruik van sinks aan wanneer u het volgende moet doen:

  • Bouw operationele use cases met lage latentie, zoals fraudedetectie, realtime analyses of aanbevelingen van klanten, waarbij gegevens naar een berichtenbus moeten stromen in plaats van naar cloudopslag. Voor workloads die een latentie van milliseconden vereisen, zie De real-timemodus gebruiken in Lakeflow-pijplijnen.
  • Schrijf getransformeerde gegevens naar tabellen die worden beheerd door een extern Delta-exemplaar, waaronder beheerde en externe tabellen in Unity Catalog.
  • Voer omgekeerde ETL uit in externe systemen, zoals het terugschrijven van verwerkte gegevens naar Apache Kafka-onderwerpen voor verbruik buiten Azure Databricks.
  • Schrijven naar een formaat dat niet standaard wordt ondersteund door Azure Databricks, met behulp van aangepaste Python-gegevensbronnen.

Sinktypen

Pijplijnen ondersteunen de volgende sinktypen:

Soort spoelbak Description
Delta-tabeldoelen Schrijf naar beheerde of externe Delta-tabellen in Unity Catalog. Geef een bestandspad of een volledig gekwalificeerde tabelnaam op.
Apache Kafka-sinks Schrijf naar Apache Kafka-topics met behulp van de Kafka-connector die in de pipeline-runtime is opgenomen.
Azure Event Hubs-doeleinden Schrijf naar Azure Event Hubs met behulp van de Kafka-interface. Maakt gebruik van dezelfde opties als Kafka-sinks.
Aangepaste Python-sinks Schrijf naar elke gegevensopslag met behulp van een aangepaste Python-gegevensbron die is geregistreerd bij spark.dataSource.register.
ForEachBatch sinks Pas aangepaste Python logica toe op elke microbatch met streaminggegevens. Gebruik dit als u naar meerdere bestemmingen moet schrijven, upserts moet uitvoeren of doelsystemen moet gebruiken die streaming-schrijfbewerkingen niet standaard ondersteunen.

Sink-API’s

Pijplijnen bieden twee API's voor het maken van sinks:

Beide typen sinks worden aangeduid als de target van een append_flow.

Limitations

  • Sinks zijn alleen beschikbaar in Python. SQL wordt niet ondersteund.
  • Alleen streamingquery's worden ondersteund. Batchquery's worden niet ondersteund.
  • Alleen append_flow kan schrijven naar sinks; create_auto_cdc_flow en andere stroomtypen worden niet ondersteund.
  • Pijplijnwachtingen worden niet ondersteund voor sinks.
  • Wanneer u een volledige verversing uitvoert, worden eerder geschreven gegevens in doelsystemen niet opgeruimd.

Aanvullende bronnen