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.
A speciális várakozási minták több adathalmazra vonatkozó elvárásokat kombinálnak az adatminőség nagy léptékű kényszerítése érdekében. Ezek a minták feltételezik, hogy tisztában van a materializált nézetek, streamtáblák és elvárások szintaxisával és szemantikájával.
Az elvárások viselkedésének és szintaxisának alapszintű áttekintéséért tekintse meg az adatminőség kezelése folyamatokkal kapcsolatos elvárásokkal című témakört.
Hordozható és újrafelhasználható elvárások
A Databricks a következő ajánlott eljárásokat javasolja a hordozhatóság javítása és a karbantartási terhek csökkentése érdekében támasztott elvárások megvalósítása során:
| Recommendation | Hatás |
|---|---|
| A várakozási definíciókat a folyamatlogikától elkülönítve tárolja. | Egyszerűen alkalmazhat elvárásokat több adatkészletre vagy folyamatra. A pipeline forráskódjának módosítása nélkül frissítheti, auditálhatja és fenntarthatja az elvárásokat. |
| Adjon hozzá egyéni címkéket a kapcsolódó elvárások csoportjainak létrehozásához. | A címkék alapján szűrheti az elvárásokat. |
| Az elvárásokat következetesen alkalmazza a hasonló adathalmazokra. | Azonos logika kiértékeléséhez használja ugyanazokat az elvárásokat több adathalmaz és folyamat között. |
Az alábbi példák bemutatják, hogyan használható egy deltatábla vagy szótár egy központi elvárások adattár létrehozásához. Az egyéni Python-függvények ezután az alábbi elvárásokat alkalmazzák egy példafolyamat adathalmazaira:
Note
Az SQL nem támogatja a fájlokból származó elvárások dinamikus betöltését.
Delta-tábla
Az alábbi példa létrehoz egy táblát a rules szabályok fenntartásához:
CREATE OR REPLACE TABLE
rules
AS SELECT
col1 AS name,
col2 AS constraint,
col3 AS tag
FROM (
VALUES
("website_not_null","Website IS NOT NULL","validity"),
("fresh_data","to_date(updateTime,'M/d/yyyy h:m:s a') > '2010-01-01'","maintained"),
("social_media_access","NOT(Facebook IS NULL AND Twitter IS NULL AND Youtube IS NULL)","maintained")
)
Az alábbi Python-példa a tábla szabályai alapján határozza meg az adatminőséggel rules kapcsolatos elvárásokat. A get_rules() függvény beolvassa a szabályokat a rules táblából, és visszaad egy Python-szótárt, amely a függvénynek átadott argumentumnak megfelelő tag szabályokat tartalmazza.
Ebben a példában a rendszer a szótárt @dp.expect_all_or_drop() dekorátorokkal alkalmazza az adatminőségi korlátozások érvényesítése érdekében.
Például a rendszer elveti a táblából azokat a rekordokat, amelyek nem összhangban vannak a validity szabályokkalraw_farmers_market:
from pyspark import pipelines as dp
from pyspark.sql.functions import expr, col
def get_rules(tag):
"""
loads data quality rules from a table
:param tag: tag to match
:return: dictionary of rules that matched the tag
"""
df = spark.read.table("rules").filter(col("tag") == tag).collect()
return {
row['name']: row['constraint']
for row in df
}
@dp.table
@dp.expect_all_or_drop(get_rules('validity'))
def raw_farmers_market():
return (
spark.read.format('csv').option("header", "true")
.load('/databricks-datasets/data.gov/farmers_markets_geographic_data/data-001/')
)
@dp.table
@dp.expect_all_or_drop(get_rules('maintained'))
def organic_farmers_market():
return (
spark.read.table("raw_farmers_market")
.filter(expr("Organic = 'Y'"))
)
Python-modul
Az alábbi példa létrehoz egy Python-modult a szabályok fenntartásához. Ebben a példában a kódot a folyamat forráskódjaként használt jegyzetfüzetmel megegyező mappában lévő rules_module.py fájlban tárolja:
def get_rules_as_list_of_dict():
return [
{
"name": "website_not_null",
"constraint": "Website IS NOT NULL",
"tag": "validity"
},
{
"name": "fresh_data",
"constraint": "to_date(updateTime,'M/d/yyyy h:m:s a') > '2010-01-01'",
"tag": "maintained"
},
{
"name": "social_media_access",
"constraint": "NOT(Facebook IS NULL AND Twitter IS NULL AND Youtube IS NULL)",
"tag": "maintained"
}
]
Az alábbi Python-példa a fájlban meghatározott szabályok alapján határozza meg az adatminőséggel rules_module.py kapcsolatos elvárásokat. A get_rules() függvény egy Python-szótárt ad vissza, amely a tag neki átadott argumentumnak megfelelő szabályokat tartalmazza.
Ebben a példában a rendszer a szótárt @dp.expect_all_or_drop() dekorátorokkal alkalmazza az adatminőségi korlátozások érvényesítése érdekében.
Például a rendszer elveti a táblából azokat a rekordokat, amelyek nem összhangban vannak a validity szabályokkalraw_farmers_market:
from pyspark import pipelines as dp
from rules_module import *
from pyspark.sql.functions import expr, col
def get_rules(tag):
"""
loads data quality rules from a table
:param tag: tag to match
:return: dictionary of rules that matched the tag
"""
return {
row['name']: row['constraint']
for row in get_rules_as_list_of_dict()
if row['tag'] == tag
}
@dp.table
@dp.expect_all_or_drop(get_rules('validity'))
def raw_farmers_market():
return (
spark.read.format('csv').option("header", "true")
.load('/databricks-datasets/data.gov/farmers_markets_geographic_data/data-001/')
)
@dp.table
@dp.expect_all_or_drop(get_rules('maintained'))
def organic_farmers_market():
return (
spark.read.table("raw_farmers_market")
.filter(expr("Organic = 'Y'"))
)
Érvényesítési táblák és folyamatvezérlési folyamat
A szakasz egyes mintái, például a sorok számának ellenőrzése és az elsődleges kulcs egyedisége külön adathalmazt, egy érvényesítési táblát határoznak meg, amely egy tulajdonságot ellenőriz más táblákban, és a problémák felszínre hozására szolgál expect_or_fail . Mielőtt érvényesítési táblázatot használna egy folyamat vezérlésére, fontos megérteni, hogy az elvárások mit szabályozhatnak, és mit nem:
- Az elvárások az adatminőséget kényszerítik ki, nem a vezénylést. A folyamaton belül az elvárások határozzák meg, hogy mely rekordok érik el a céladatkészletet:
warnmegőrzi az érvénytelen rekordokat és a rekordmetrikákat,dropeltávolítja őket, ésfailleállítja a jogsértő folyamatot. A cél annak biztosítása, hogy csak tiszta adatok haladjanak át, ne feltételesen futtassa vagy hagyja ki a folyamat más részeit. -
expect_or_failviselkedése a folyamat futási módjától függ. Egy aktivált folyamatban egy sikertelen várakozás meghiúsul, és csak a folyamat frissítését veti vissza; az ugyanabban a folyamatban lévő többi folyamat továbbra is egymástól függetlenül frissül. Egy folyamatos folyamatban a sikertelen várakozás leállítja a folyamatot és az összes függő folyamatot. Lásd: Sikertelen érvénytelen rekordok esetén. - Az érvényesítési táblák nem nyitják meg az alsóbb rétegbeli táblákat. Ha egy érvényesítési táblát egy másik adatkészletből olvas be, az adathalmaz nem várja meg az érvényesítési eredményt, így a sikertelen érvényesítés nem akadályozza meg az alsóbb rétegbeli táblák frissítését.
Ha le szeretné állítani az alsóbb rétegbeli feldolgozást, ha az ellenőrzés sikertelen, ossza fel az érvényesítési logikát és az alsóbb rétegbeli munkát külön folyamatokra, és vezényelje őket egy feladattal, így az alsóbb rétegbeli folyamat feladata az érvényesítési folyamat tevékenységétől függ. Mivel egy folyamattevékenység meghiúsul a frissítés sikertelensége esetén, az alsóbb rétegbeli tevékenység csak akkor fut, ha az érvényesítési folyamat sikeres. Általánosabban, ha feltételes végrehajtásra vagy összetett függőségekre van szüksége, koordináljon több folyamatot egy feladattal ahelyett, hogy a logikát egyetlen folyamatba építi. Lásd: Folyamatok futtatása munkafolyamatban.
Sorok számának érvényesítése
Az alábbi példa a sorok számának egyenlőségét ellenőrzi table_a és table_b között annak ellenőrzésére, hogy az átalakítások során nem vesztek-e el adatok.
Python
@dp.materialized_view(
name="count_verification",
comment="Validates equal row counts between tables"
)
@dp.expect_or_fail("no_rows_dropped", "a_count == b_count")
def validate_row_counts():
return spark.sql("""
SELECT * FROM
(SELECT COUNT(*) AS a_count FROM table_a),
(SELECT COUNT(*) AS b_count FROM table_b)""")
SQL
CREATE OR REFRESH MATERIALIZED VIEW count_verification(
CONSTRAINT no_rows_dropped EXPECT (a_count == b_count)
) AS SELECT * FROM
(SELECT COUNT(*) AS a_count FROM table_a),
(SELECT COUNT(*) AS b_count FROM table_b)
Hiányzó rekordészlelés
Az alábbi példa ellenőrzi, hogy az összes várt rekord szerepel-e a report táblában:
Python
@dp.materialized_view(
name="report_compare_tests",
comment="Validates no records are missing after joining"
)
@dp.expect_or_fail("no_missing_records", "r_key IS NOT NULL")
def validate_report_completeness():
return (
spark.read.table("validation_copy").alias("v")
.join(
spark.read.table("report").alias("r"),
on="key",
how="left_outer"
)
.select(
"v.*",
"r.key as r_key"
)
)
SQL
CREATE OR REFRESH MATERIALIZED VIEW report_compare_tests(
CONSTRAINT no_missing_records EXPECT (r_key IS NOT NULL)
)
AS SELECT v.*, r.key as r_key FROM validation_copy v
LEFT OUTER JOIN report r ON v.key = r.key
Elsődleges kulcs egyedisége
Az alábbi példa az elsődleges kulcsokra vonatkozó korlátozásokat ellenőrzi a táblákon:
Python
@dp.materialized_view(
name="report_pk_tests",
comment="Validates primary key uniqueness"
)
@dp.expect_or_fail("unique_pk", "num_entries = 1")
def validate_pk_uniqueness():
return (
spark.read.table("report")
.groupBy("pk")
.count()
.withColumnRenamed("count", "num_entries")
)
SQL
CREATE OR REFRESH MATERIALIZED VIEW report_pk_tests(
CONSTRAINT unique_pk EXPECT (num_entries = 1)
)
AS SELECT pk, count(*) as num_entries
FROM report
GROUP BY pk
Sémafejlődési minta
Az alábbi példa bemutatja, hogyan kezelheti a sémafejlődést a további oszlopok esetében. Ezt a mintát akkor használja, ha adatforrásokat migrál vagy több verziójú felsőbb rétegbeli adatot kezel, biztosítva a visszamenőleges kompatibilitást az adatminőség kényszerítése mellett:
Python
@dp.table
@dp.expect_all_or_fail({
"required_columns": "col1 IS NOT NULL AND col2 IS NOT NULL",
"valid_col3": "CASE WHEN col3 IS NOT NULL THEN col3 > 0 ELSE TRUE END"
})
def evolving_table():
# Legacy data (V1 schema)
legacy_data = spark.read.table("legacy_source")
# New data (V2 schema)
new_data = spark.read.table("new_source")
# Combine both sources
return legacy_data.unionByName(new_data, allowMissingColumns=True)
SQL
CREATE OR REFRESH MATERIALIZED VIEW evolving_table(
-- Merging multiple constraints into one as expect_all is Python-specific API
CONSTRAINT valid_migrated_data EXPECT (
(col1 IS NOT NULL AND col2 IS NOT NULL) AND (CASE WHEN col3 IS NOT NULL THEN col3 > 0 ELSE TRUE END)
) ON VIOLATION FAIL UPDATE
) AS
SELECT * FROM new_source
UNION
SELECT *, NULL as col3 FROM legacy_source;
Tartományalapú érvényesítési minta
Az alábbi példa bemutatja, hogyan érvényesítheti az új adatpontokat az előzménystatisztikai tartományokkal szemben, és hogyan azonosíthatja az ön adatfolyamának kiugró értékeit és anomáliáit.
Python
@dp.view
def stats_validation_view():
# Calculate statistical bounds from historical data
bounds = spark.sql("""
SELECT
avg(amount) - 3 * stddev(amount) as lower_bound,
avg(amount) + 3 * stddev(amount) as upper_bound
FROM historical_stats
WHERE
date >= CURRENT_DATE() - INTERVAL 30 DAYS
""")
# Join with new data and apply bounds
return spark.read.table("new_data").crossJoin(bounds)
@dp.table
@dp.expect_or_drop(
"within_statistical_range",
"amount BETWEEN lower_bound AND upper_bound"
)
def validated_amounts():
return spark.read.table("stats_validation_view")
SQL
CREATE OR REFRESH MATERIALIZED VIEW stats_validation_view AS
WITH bounds AS (
SELECT
avg(amount) - 3 * stddev(amount) as lower_bound,
avg(amount) + 3 * stddev(amount) as upper_bound
FROM historical_stats
WHERE date >= CURRENT_DATE() - INTERVAL 30 DAYS
)
SELECT
new_data.*,
bounds.*
FROM new_data
CROSS JOIN bounds;
CREATE OR REFRESH MATERIALIZED VIEW validated_amounts (
CONSTRAINT within_statistical_range EXPECT (amount BETWEEN lower_bound AND upper_bound)
)
AS SELECT * FROM stats_validation_view;
Érvénytelen rekordok karanténba helyezése
Ez a minta egyesíti az elvárásokat ideiglenes táblákkal és nézetekkel az adatminőségi metrikák nyomon követéséhez a folyamatfrissítések során, és lehetővé teszi a különálló feldolgozási útvonalakat az alsóbb rétegbeli műveletek érvényes és érvénytelen rekordjaihoz.
Python
from pyspark import pipelines as dp
from pyspark.sql.functions import expr
rules = {
"valid_pickup_zip": "(pickup_zip IS NOT NULL)",
"valid_dropoff_zip": "(dropoff_zip IS NOT NULL)",
}
quarantine_rules = "NOT({0})".format(" AND ".join(rules.values()))
@dp.view
def raw_trips_data():
return spark.readStream.table("samples.nyctaxi.trips")
@dp.table(
temporary=True,
partition_cols=["is_quarantined"],
)
@dp.expect_all(rules)
def trips_data_quarantine():
return (
spark.readStream.table("raw_trips_data").withColumn("is_quarantined", expr(quarantine_rules))
)
@dp.view
def valid_trips_data():
return spark.read.table("trips_data_quarantine").filter("is_quarantined=false")
@dp.view
def invalid_trips_data():
return spark.read.table("trips_data_quarantine").filter("is_quarantined=true")
SQL
CREATE TEMPORARY STREAMING LIVE VIEW raw_trips_data AS
SELECT * FROM STREAM(samples.nyctaxi.trips);
CREATE OR REFRESH TEMPORARY STREAMING TABLE trips_data_quarantine(
-- Option 1 - merge all expectations to have a single name in the pipeline event log
CONSTRAINT quarantined_row EXPECT (pickup_zip IS NOT NULL OR dropoff_zip IS NOT NULL),
-- Option 2 - Keep the expectations separate, resulting in multiple entries under different names
CONSTRAINT invalid_pickup_zip EXPECT (pickup_zip IS NOT NULL),
CONSTRAINT invalid_dropoff_zip EXPECT (dropoff_zip IS NOT NULL)
)
PARTITIONED BY (is_quarantined)
AS
SELECT
*,
NOT ((pickup_zip IS NOT NULL) and (dropoff_zip IS NOT NULL)) as is_quarantined
FROM STREAM(raw_trips_data);
CREATE TEMPORARY LIVE VIEW valid_trips_data AS
SELECT * FROM trips_data_quarantine WHERE is_quarantined=FALSE;
CREATE TEMPORARY LIVE VIEW invalid_trips_data AS
SELECT * FROM trips_data_quarantine WHERE is_quarantined=TRUE;