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.
Egy folyamatlánc több, egymáshoz szinte teljesen azonos folyamatot is tartalmazhat, amelyek csak néhány paraméterben térnek el. Ezeknek a folyamatoknak a explicit meghatározása hibalehetőséget jelent, redundáns és nehezen karbantartható. A Python belső függvényekkel történő metaprogramozás dinamikusan hoz létre ismétlődő folyamatokat, és mindegyik meghívás különböző paramétereket biztosít.
Metaprogramozási áttekintés
A Lakeflow-folyamatok metaprogramozása Python belső függvényeket használ. Mivel ezeket a függvényeket a csővezeték futási környezete lustán kiértékeli, becsomagolhatja a @dp.table dekorátorokat egy gyári függvénybe, és többször meghívhatja a gyárat különböző paraméterekkel. Minden hívás regisztrál egy új folyamatot kód duplikálása nélkül.
A for hurkok Lakeflow-folyamatokkal való használatával kapcsolatos részletekért lásd a Táblák létrehozása egy for ciklusban című részt.
Példa: tűzoltóság válaszideje
Az alábbi példa az alapértelmezett tűzoltósági adatkészletet használja, hogy megkeresse azokat a környékeket, amelyek a leggyorsabb vészhelyzeti válaszidőkkel rendelkeznek az egyes hívástípusok esetén. Metaprogramozás nélkül szinte azonos tábladefiníciókat kell írnia minden hívástípushoz (Riasztások, Struktúratűz, Orvosi incidens). A metaprogramozással egyetlen gyári függvény hozza létre az összeset.
1. lépés: A nyers feldolgozási tábla meghatározása
import functools
from pyspark import pipelines as dp
from pyspark.sql.functions import *
@dp.table(
name="raw_fire_department",
comment="raw table for fire department response"
)
@dp.expect_or_drop("valid_received", "received IS NOT NULL")
@dp.expect_or_drop("valid_response", "responded IS NOT NULL")
@dp.expect_or_drop("valid_neighborhood", "neighborhood != 'None'")
def get_raw_fire_department():
return (
spark.read.format('csv')
.option('header', 'true')
.option('multiline', 'true')
.load('/databricks-datasets/timeseries/Fires/Fire_Department_Calls_for_Service.csv')
.withColumnRenamed('Call Type', 'call_type')
.withColumnRenamed('Received DtTm', 'received')
.withColumnRenamed('Response DtTm', 'responded')
.withColumnRenamed('Neighborhooods - Analysis Boundaries', 'neighborhood')
.select('call_type', 'received', 'responded', 'neighborhood')
)
2. lépés: A flow factory függvény definiálása
A generate_tables gyári függvény két táblát regisztrál minden hívástípushoz: egy szűrt hívástáblát és egy rangsorolt válaszidő-táblát. Mindkettő belső függvényként jön létre és van dekorálva a @dp.table-nal.
all_tables = []
def generate_tables(call_table, response_table, filter):
@dp.table(
name=call_table,
comment="top level tables by call type"
)
def create_call_table():
return spark.sql("""
SELECT
unix_timestamp(received,'M/d/yyyy h:m:s a') as ts_received,
unix_timestamp(responded,'M/d/yyyy h:m:s a') as ts_responded,
neighborhood
FROM raw_fire_department
WHERE call_type = '{filter}'
""".format(filter=filter))
@dp.table(
name=response_table,
comment="top 10 neighborhoods with fastest response time"
)
def create_response_table():
return spark.sql("""
SELECT
neighborhood,
AVG((ts_received - ts_responded)) as response_time
FROM {call_table}
GROUP BY 1
ORDER BY response_time
LIMIT 10
""".format(call_table=call_table))
all_tables.append(response_table)
3. lépés: A gyár meghívása és az összefoglaló tábla definiálása
Minden hívástípushoz egyszer hívja meg a gyárat, majd adjon meg egy összegző táblát, amely egyesíteni tudja az eredményeket, hogy megtalálja az összes kategóriában leggyakrabban megjelenő környékeket.
generate_tables("alarms_table", "alarms_response", "Alarms")
generate_tables("fire_table", "fire_response", "Structure Fire")
generate_tables("medical_table", "medical_response", "Medical Incident")
@dp.table(
name="best_neighborhoods",
comment="which neighbor appears in the best response time list the most"
)
def summary():
target_tables = [dp.read(t) for t in all_tables]
unioned = functools.reduce(lambda x, y: x.union(y), target_tables)
return (
unioned.groupBy(col("neighborhood"))
.agg(count("*").alias("score"))
.orderBy(desc("score"))
)
A folyamat futtatása után hasonló táblákat hoz létre, például a következő gráfot:
főbb fogalmak
-
A belső függvények lazán vannak regisztrálva: A
@dp.tabledekorátor nem futtatja azonnal a függvényt. Regisztrálja a függvényt a folyamat futtatókörnyezetével, amely a végrehajtás megkezdése előtt feloldja a teljes adatfolyam-gráfot. -
A lezárások rögzítik a paramétereket: Minden belső függvény lezárja a gyárnak átadott paramétereket (
call_table,response_table,filter), így minden regisztrált folyamat a saját elkülönített értékkészletét használja. -
Dinamikus táblázatlisták: A programozott módon létrehozott táblanevek nyomon követéséhez hasonló
all_tableslista használatával egyszerűen hivatkozhat rájuk később (például egy egyesítésben vagy illesztésben).