Använda sänkor i pipelines

Använd Lakeflow-pipeline-sinkAPI:t med flöden för att skriva poster som har transformerats av en pipeline till en extern datamottagare. Externa datamottagare inkluderar hanterade och externa tabeller i Unity Catalog och händelseströmningstjänster som Apache Kafka eller Azure Event Hubs. Du kan också använda datamottagare för att skriva till anpassade datakällor genom att skriva Python-kod för datakällan.

En översikt över mottagarbegrepp och när du ska använda dem finns i Sinks in Lakeflow pipelines (Mottagare i Lakeflow-pipelines).

Anmärkning

Arbetsflöde för mottagarnod

När händelsedata matas in från en strömmande källa till din pipeline bearbetar och förfinar du dessa data i transformeringar i din pipeline. Sedan använder du append-flödesbearbetning för att strömma de transformerade dataposterna till en slutpunkt. Du skapar den här sinken med hjälp av funktionen create_sink(). För mer information om create_sink-funktionen, se mottagare-API-referensen .

Om du har en pipeline som skapar eller bearbetar dina strömmande händelsedata och förbereder dataposter för skrivning, är du redo att använda en sink.

Implementeringen av en mottagare består av två steg:

  1. Skapa mottagaren.
  2. Använd ett tilläggsflöde eller uppdateringsflöde för att skriva de förberedda posterna till mottagaren.

Skapa en mottagare

Databricks har stöd för flera typer av destinationspunkter där du skriver dina poster som bearbetats från din dataström.

  • Deltatabellsänkor (inklusive hanterade och externa Unity-katalogtabeller)
  • Apache Kafka-sänkor
  • Azure Event Hubs-sänkor
  • Anpassade mottagare skrivna i Python med hjälp av anpassade Python-datakällor

Nedan visas exempel på konfigurationer för Delta-, Kafka- och Azure Event Hubs-mottagare och python-anpassade datakällor:

Deltamottagare

Så här skapar du en Delta-sänka via filsökväg:

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

Så här skapar du en Delta-sink utifrån tabellnamn med hjälp av en fullständigt kvalificerad katalog- och schemasökväg.

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

Kafka- och Azure Event Hubs-utsläpp

Den här koden fungerar för både Apache Kafka- och Azure Event Hubs-sänkor.

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
  }
)

credential_name är en referens till en autentiseringsuppgift för Unity Catalog-tjänsten. Mer information finns i Använda autentiseringsuppgifter för Unity Catalog-tjänsten för att ansluta till externa molntjänster.

Anpassade Python-datakällor

Förutsatt att du har en anpassad Python-datakälla registrerad som my_custom_datasourcekan följande kod skriva till datakällan.

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")

Mer information om hur du skapar anpassade datakällor i Python finns i PySpark-anpassade datakällor.

För mer information om funktionen create_sink, se API-referensen för sink .

När mottagaren har skapats kan du börja strömma bearbetade poster till mottagaren.

Skriva till en mottagare med ett tilläggsflöde

När diskhonen har skapats är nästa steg att skriva bearbetade poster till den genom att definiera den som målet för poster som matas ut av ett tilläggsflöde. Du gör detta genom att ange ditt handfat som värdet target i dekoratören append_flow.

  • För hanterade och externa unity-katalogtabeller använder du formatet delta och anger sökvägen eller tabellnamnet i alternativen. Din pipeline måste konfigureras för att använda Unity Catalog.
  • För Apache Kafka-ämnen använder du formatet kafka och anger ämnesnamn, anslutningsinformation och autentiseringsinformation i alternativen. Det här är samma alternativ som en Spark Structured Streaming Kafka-mottagare stöder. Se Konfigurera Kafka Structured Streaming-skrivaren.
  • För Azure Event Hubs använder du formatet kafka och anger händelsehubbarnas namn, anslutningsinformation och autentiseringsinformation i alternativen. Det här är samma alternativ som stöds i en Spark Structured Streaming Event Hubs-mottagare som använder Kafka-gränssnittet. Se Autentisering.

Nedan visas exempel på hur du konfigurerar flöden för att skriva till Delta-, Kafka- och Azure Event Hubs-mottagare med poster som bearbetas av din pipeline.

Delta-handfat

@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- och Azure Event Hubs-utsläpp

@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")
)

Parametern value är obligatorisk för en Azure Event Hubs-mottagare. Ytterligare parametrar som key, partition, headersoch topic är valfria.

Mer information om dekoratören finns i append_flowStandardflöden och tilläggsflöden.

Begränsningar

  • Endast Python-API:et stöds. SQL stöds inte.

  • Endast strömmande frågeanrop stöds. Batchfrågor stöds inte.

  • Endast append_flow och update_flow kan användas för att skriva till mottagare. Andra flöden, till exempel create_auto_cdc_flow, stöds inte och du kan inte använda en sink i en definition av pipelinedatauppsättning. Följande stöds till exempel inte:

    @table("from_sink_table")
    def fromSink():
      return read_stream("my_sink")
    
  • För Delta-sinkar måste tabellnamnet vara fullständigt kvalificerat. För Unity Catalog hanterade externa tabeller måste tabellnamnet vara i formatet <catalog>.<schema>.<table>. För Hive-metastore måste det vara i formatet <schema>.<table>.

  • När du kör en fullständig uppdatering tas inte tidigare beräknade resultatdata bort i datakällorna. Det innebär att all ombearbetad data läggs till i sänkan och att befintlig data inte ändras.

  • Pipelineförväntningar stöds inte.

  • Serverlös kontroll av utgående trafik har endast stöd för Kafka- och Delta Lake-sink-connectorer. Se Vad är serverlös hantering av utgående trafik?.

Ytterligare resurser