Sinks gebruiken in pijplijnen

Gebruik de Lakeflow-pijplijn-API sink met stromen om records te schrijven die zijn getransformeerd door een pijplijn naar een externe gegevenssink. Externe gegevenssinks zijn beheerde en externe tabellen van Unity Catalog en streamingservices voor gebeurtenissen, zoals Apache Kafka of Azure Event Hubs. U kunt ook gegevenssinks gebruiken om naar aangepaste gegevensbronnen te schrijven door Python-code voor die gegevensbron te schrijven.

Zie Sinks in Lakeflow-pijplijnen voor een overzicht van sinkconcepten en wanneer ze moeten worden gebruikt.

Opmerking

Sink-werkstroom

Wanneer gebeurtenisgegevens worden opgenomen uit een streamingbron in uw pijplijn, verwerkt en verfijnt u deze gegevens in transformaties in uw pijplijn. Vervolgens gebruikt u de verwerking van toevoegstromen om de getransformeerde gegevensrecords naar een sink te streamen. U maakt deze sink met behulp van de functie create_sink(). Zie de create_sink voor meer informatie over de functie.

Als u een pijplijn hebt die uw streaming-gebeurtenisgegevens maakt of verwerkt en gegevensrecords voorbereidt voor schrijven, kunt u een sink gebruiken.

Het implementeren van een sink bestaat uit twee stappen:

  1. Maak de sink.
  2. Gebruik een toevoegstroom of updatestroom om de voorbereide records naar de sink te schrijven.

Een sink maken

Databricks ondersteunt verschillende typen doelsinks waarin u uw records schrijft die zijn verwerkt vanuit uw streamgegevens:

  • Delta-tabelsinks (inclusief beheerde en externe tabellen in Unity Catalog)
  • Apache Kafka-sinks
  • Azure Event Hubs-doeleinden
  • Aangepaste sinks geschreven in Python, gebruikmakend van de Python-gegevensbronnen op maat

Hieronder ziet u voorbeelden van configuraties voor Delta-, Kafka- en Azure Event Hubs-sinks en aangepaste Python-gegevensbronnen:

Delta gootstenen

Een Delta-sink maken via een bestandspad:

dp.create_sink(
  name = "delta_sink",
  format = "delta",
  options = {"path": "/Volumes/catalog_name/schema_name/volume_name/path/to/data"}
)

Een Delta-sink maken op tabelnaam met behulp van een volledig gekwalificeerde catalogus en schemapad:

dp.create_sink(
  name = "delta_sink",
  format = "delta",
  options = { "tableName": "catalog_name.schema_name.table_name" }
)

Kafka- en Azure Event Hubs-sinks

Deze code werkt voor zowel Apache Kafka als Azure Event Hubs-sinks.

credential_name = "<service-credential>"
eh_namespace_name = "dp-eventhub"
bootstrap_servers = f"{eh_namespace_name}.servicebus.windows.net:9093"
topic_name = "dp-sink"

dp.create_sink(
name = "eh_sink",
format = "kafka",
options = {
    "databricks.serviceCredential": credential_name,
    "kafka.bootstrap.servers": bootstrap_servers,
    "topic": topic_name
  }
)

Dit credential_name is een verwijzing naar een Unity Catalog-servicereferentie. Voor meer informatie, zie Unity Catalog-service referenties gebruiken om te verbinden met externe cloudservices.

Aangepaste Python-gegevensbronnen

Ervan uitgaande dat u een aangepaste Python-gegevensbron hebt geregistreerd als my_custom_datasource, kan de volgende code naar die gegevensbron schrijven.

from pyspark import pipelines as dp

# Assume `my_custom_datasource` is a custom Python streaming
# data source that writes data to your system.

# Create Lakeflow pipelines sink using my_custom_datasource
dp.create_sink(
    name="custom_sink",
    format="my_custom_datasource",
    options={
        <options-needed-for-custom-datasource>
    }
)

# Create append flow to send data to RequestBin
@dp.append_flow(name="flow_to_custom_sink", target="custom_sink")
def flow_to_custom_sink():
    return read_stream("my_source_data")

Zie Aangepaste gegevensbronnen van PySpark voor meer informatie over het maken van aangepaste gegevensbronnen in Python.

Zie de create_sinkvoor meer informatie over het gebruik van de functie.

Nadat uw sink is gemaakt, kunt u beginnen met het streamen van verwerkte records naar de sink.

Schrijven naar een sink met een toevoegstroom

Nadat uw sink is aangemaakt, is de volgende stap om verwerkte records naar de sink te schrijven door deze op te geven als het doel voor de uitvoer van records via een append-flow. U doet dit door uw sink op te geven als de target waarde in de append_flow decorator.

  • Voor beheerde en externe tabellen van Unity Catalog gebruikt u de indeling delta en geeft u het pad of de tabelnaam op in opties. Uw pijplijn moet zijn geconfigureerd voor het gebruik van Unity Catalog.
  • Voor Apache Kafka-onderwerpen gebruikt u de indeling kafka en geeft u de onderwerpnaam, verbindingsinformatie en verificatiegegevens op in de opties. Dit zijn dezelfde opties die een Spark Structured Streaming Kafka-sink ondersteunt. Zie De Kafka Structured Streaming Writer configureren.
  • Voor Azure Event Hubs gebruikt u de indeling kafka en geeft u de naam, verbindingsgegevens en verificatiegegevens van Event Hubs op in de opties. Dit zijn dezelfde opties die worden ondersteund in een Spark Structured Streaming Event Hubs-sink die gebruikmaakt van de Kafka-interface. Zie Verificatie.

Hieronder ziet u voorbeelden van het instellen van stromen voor het schrijven naar Delta-, Kafka- en Azure Event Hubs-sinks met records die door uw pijplijn worden verwerkt.

Delta-wastafel

@dp.append_flow(name = "delta_sink_flow", target="delta_sink")
def delta_sink_flow():
  return(
  spark.readStream.table("spark_referrers")
  .selectExpr("current_page_id", "referrer", "current_page_title", "click_count")
)

Kafka- en Azure Event Hubs-sinks

@dp.append_flow(name = "kafka_sink_flow", target = "eh_sink")
def kafka_sink_flow():
return (
  spark.readStream.table("spark_referrers")
  .selectExpr("cast(current_page_id as string) as key", "to_json(struct(referrer, current_page_title, click_count)) AS value")
)

De parameter value is verplicht voor een Azure Event Hubs-sink. Aanvullende parameters, zoals key, partition, headersen topic zijn optioneel.

Zie append_flow voor meer informatie over de decorator.

Beperkingen

  • Alleen de Python-API wordt ondersteund. SQL wordt niet ondersteund.

  • Alleen streamingquery's worden ondersteund. Batchquery's worden niet ondersteund.

  • Alleen append_flow en update_flow kan worden gebruikt om naar sinks te schrijven. Andere stromen, zoals create_auto_cdc_flow, worden niet ondersteund en u kunt geen sink gebruiken in een definitie van een pijplijngegevensset. Het volgende wordt bijvoorbeeld niet ondersteund:

    @table("from_sink_table")
    def fromSink():
      return read_stream("my_sink")
    
  • Voor Delta-sinks moet de tabelnaam volledig gespecificeerd zijn. Met name voor door Unity Catalog beheerde externe tabellen moet de tabelnaam van het formulier <catalog>.<schema>.<table>zijn. Voor de Hive-metastore moet deze in de vorm <schema>.<table>zijn.

  • Het uitvoeren van een volledige vernieuwingsupdate ruimt niet automatisch eerder berekende gegevens op in de sinks. Dit betekent dat alle opnieuw verwerkte gegevens worden toegevoegd aan de sink en dat bestaande gegevens niet worden gewijzigd.

  • De verwachtingen van pijplijnen worden niet ondersteund.

  • Serverloze uitgaand-verkeerbeheer ondersteunt alleen Kafka- en Delta Lake-sink-connectors. Zie Wat is serverloos egressbeheer?.

Aanvullende bronnen