Python gebruiken met zelfstandige pijplijnen

U kunt zelfstandige gerealiseerde weergaven en streamingtabellen maken en vernieuwen vanuit een notebook met behulp van Python. Ontwerp uw pijplijn in een Python notebook en voer deze uit met spark.sql(). Hiermee kunt u zelfstandige pijplijnen naast uw andere Python notebookwerkstromen beheren.

Python-broncode voor zelfstandige pijplijnen vereist een notebook dat is gekoppeld aan serverloze algemene rekenkracht. U kunt Python niet gebruiken om zelfstandige pijplijnen te maken of te vernieuwen vanuit een Databricks SQL-warehouse, omdat een warehouse SQL-instructies uitvoert, niet Python notebooks. Als u in plaats daarvan een SQL-warehouse wilt gebruiken, zie Zelfstandige gematerialiseerde weergaven gebruiken en Zelfstandige streamingtabellen gebruiken.

Important

Het maken en vernieuwen van zelfstandige materialized views en streamingtabellen vanuit een notebook op serverloze compute voor algemene doeleinden is bèta en beschikbaar in geselecteerde regio's. Zie Notebooks.

Requirements

Als u zelfstandige pijplijnen wilt maken en vernieuwen met Python, hebt u een notebook nodig dat is gekoppeld aan serverloze algemene berekeningen op Databricks Runtime 18.1 of hoger. Zie Notebooks voor de volledige lijst met vereisten, waaronder regionale beschikbaarheid en machtigingen.

Hoe werkt het?

Geef in een Python notebook dezelfde instructies door die u zou uitvoeren vanuit een Databricks SQL Warehouse naar spark.sql(). De syntaxis van de zelfstandige gematerialiseerde weergave en de streamingtabel is identiek; alleen de manier waarop u de opdracht indient, verschilt. Net als bij een magazijn voert elke CREATE of REFRESH instructie een serverloze pijplijn uit om de bewerking te verwerken.

De spark sessie is standaard beschikbaar in Azure Databricks notebooks, dus er is geen import vereist.

Een gerealiseerde weergave maken

In het volgende voorbeeld wordt de gerealiseerde weergave gemaakt mv1 uit de basistabel base_table1:

spark.sql("""
  CREATE OR REPLACE MATERIALIZED VIEW mv1
  AS SELECT
    date,
    sum(sales) AS sum_of_sales
  FROM base_table1
  GROUP BY date
""")

Zie CREATE MATERIALIZED VIEW voor volledige details, zoals geplande en geactiveerde vernieuwingen.

Een streamingtabel maken

In het volgende voorbeeld wordt de streamingtabel sales gemaakt op basis van de raw_data tabel:

spark.sql("""
  CREATE OR REFRESH STREAMING TABLE sales
  AS SELECT product, price FROM STREAM raw_data
""")

Zie CREATE STREAMING TABLE voor meer informatie, waaronder het laden van bestanden met automatisch laden en plannen.

Een gematerialiseerde weergave of streamingtabel vernieuwen

Gebruik een REFRESH instructie om een zelfstandige tabel bij te werken met de meest recente gegevens uit de bron:

spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")

Op serverloze algemene berekeningen zijn vernieuwingen synchroon. Asynchrone vernieuwingen (het ASYNC trefwoord) worden niet ondersteund. Zie Serverloze algemene rekenkracht.

Instructies voor parameteriseren

Als u waarden uit uw Python code wilt doorgeven aan een instructie in plaats van ze te coderen, gebruikt u benoemde parametermarkeringen in de SQL en geeft u hun waarden op via het args argument van spark.sql(). Gebruik een markering, zoals :min_sales rechtstreeks voor letterlijke waarden. Verpakt de markering alleen in IDENTIFIER() wanneer de parameter een objectnaam is, zoals een tabel, weergave of schema, omdat id's niet kunnen worden vervangen als tekenreekswaarden zonder opmaak.

In het volgende voorbeeld worden zowel de gerealiseerde weergavenaam als een filterwaarde geparameteriseerd:

mv_name = "main.sales.regional_sales"
min_sales = 1000

spark.sql("""
  CREATE OR REPLACE MATERIALIZED VIEW IDENTIFIER(:mv)
  AS SELECT
    region,
    sum(sales) AS sum_of_sales
  FROM base_table1
  WHERE sales > :min_sales
  GROUP BY region
""", args={
  "mv": mv_name,
  "min_sales": min_sales,
})

Raadpleeg parametermarkeringen en IDENTIFIER-component voor meer informatie.

Andere instructies uitvoeren

U kunt elke zelfstandige gerealiseerde weergave- of streamingtabelinstructie uitvoeren vanuit een Python notebook door deze door te geven aanspark.sql(), inclusief instructies voor het plannen van vernieuwingen, het wijzigen van een tabel of het verwijderen van een tabel. Raadpleeg Zelfstandige gematerialiseerde weergaven gebruiken en Zelfstandige streamingtabellen gebruiken voor informatie over het gebruik van gematerialiseerde weergaven en streamingtabellen, inclusief SQL-syntaxis.

Limitations

Zelfstandige gerealiseerde weergaven en streamingtabellen die zijn gemaakt op serverloze algemene berekeningen hebben extra beperkingen, zoals geen ondersteuning voor asynchrone vernieuwingen en geen kostentoewijzing per tabel. Zie Serverloze algemene berekeningen voor de volledige lijst.

Aanvullende bronnen