หมายเหตุ
การเข้าถึงหน้านี้ต้องได้รับการอนุญาต คุณสามารถลอง ลงชื่อเข้าใช้หรือเปลี่ยนไดเรกทอรีได้
การเข้าถึงหน้านี้ต้องได้รับการอนุญาต คุณสามารถลองเปลี่ยนไดเรกทอรีได้
ในสถาปัตยกรรมเหรียญ คุณต้องบังคับใช้คุณภาพของข้อมูลในทุกขั้นตอน คุณภาพข้อมูลที่ไม่ดีอาจนําไปสู่ข้อมูลเชิงลึกที่ไม่ถูกต้องและความไร้ประสิทธิภาพในการปฏิบัติงาน
บทความนี้อธิบายวิธีการใช้การตรวจสอบคุณภาพข้อมูลในมุมมองทะเลสาบที่เป็นรูปธรรม (MLV) ใน Microsoft Fabric
นําคุณภาพข้อมูลไปใช้
ในมุมมองทะเลสาบที่สร้างขึ้นจริง (MLVs) ใน Fabric ให้รักษาคุณภาพข้อมูลโดยการกําหนดข้อจํากัดในมุมมองของคุณ หากไม่มีการตรวจสอบที่ชัดเจนปัญหาข้อมูลเล็กน้อยอาจเพิ่มเวลาในการประมวลผลหรือทําให้ไปป์ไลน์ล้มเหลว
เมื่อแถวละเมิดข้อจํากัด คุณสามารถใช้การดําเนินการอย่างใดอย่างหนึ่งต่อไปนี้:
ล้มเหลว: หยุดการรีเฟรช MLV เมื่อมีการละเมิดข้อจํากัดครั้งแรก นี่คือลักษณะการทํางานเริ่มต้น แม้ว่าคุณจะไม่ได้ระบุ
FAILDROP: ประมวลผลต่อไปและลบเรกคอร์ดที่ละเมิดข้อจํากัด มุมมองสายข้อมูลแสดงจํานวนเรกคอร์ดที่ลดลง
หมายเหตุ
ถ้าคุณกําหนดทั้งการดําเนินการ DROP และ FAIL ใน MLV การดําเนินการ FAIL จะมีความสําคัญเหนือกว่า
กําหนดการตรวจสอบคุณภาพข้อมูลในมุมมองทะเลสาบที่เป็นรูปธรรม
เมื่อคุณสร้างมุมมองทะเลสาบที่เป็นรูปธรรม คุณสามารถกําหนดข้อจํากัด ซึ่งเป็นกฎคุณภาพข้อมูลที่ตรวจสอบความถูกต้องของแต่ละแถวในระหว่างการรีเฟรช ข้อจํากัดคือนิพจน์บูลีนที่ทุกแถวต้องตอบสนอง แถวที่ผ่านจะถูกเขียนลงในตารางผลลัพธ์ แถวที่ล้มเหลวจะได้รับการจัดการตามการตั้งค่าการละเมิด: แถวเหล่านี้จะถูกทิ้งอย่างเงียบ ๆ หรือทําให้การรีเฟรชทั้งหมดล้มเหลว
ตัวอย่างต่อไปนี้กําหนดข้อจํากัด cust_blankซึ่งจะตรวจสอบว่า customerName เขตข้อมูลไม่ใช่ null หรือไม่ ข้อจํากัดไม่รวมแถวที่มีค่า Null customerName จากการประมวลผล
CREATE OR REPLACE MATERIALIZED LAKE VIEW IF NOT EXISTS silver.customers_enriched
(CONSTRAINT cust_blank CHECK (customerName is not null) on MISMATCH DROP)
AS
SELECT
c.customerID,
c.customerName,
c.contact,
CASE
WHEN COUNT(o.orderID) OVER (PARTITION BY c.customerID) > 0 THEN TRUE
ELSE FALSE
END AS has_orders
FROM bronze.customers c LEFT JOIN bronze.orders o
ON c.customerID = o.customerID;
ระบบสร้างขึ้นในฟังก์ชั่น
ฟังก์ชัน Spark/SQL ในตัว เช่น UPPER(), LOWER(), TRIM(), COALESCE(), INITCAP() และ DATE_FORMAT() ได้รับการสนับสนุนอย่างเต็มที่ในบริบท MLV ทั้งหมดสําหรับทั้ง CREATE และ LINEAGE รีเฟรช
CREATE MATERIALIZED LAKE VIEW sample_lakehouse.silver.names (
CONSTRAINT substring_check
CHECK (SUBSTRING(name, 1, 2) = 'Al') ON MISMATCH drop
) AS
SELECT id, name
FROM (VALUES (1, 'Alice'), (2, 'Bob'), (3, 'Ann')) AS t(id, name)
หมายเหตุ
ฟังก์ชันของระบบเป็นตัวเลือกที่ง่ายและน่าเชื่อถือที่สุด พวกเขาไม่จําเป็นต้องลงทะเบียน ทํางานในทุกบริบท และได้รับการสนับสนุนอย่างเต็มที่ในระหว่างการรีเฟรช LINEAGE
UDF – กําหนดและลงทะเบียนในสมุดบันทึกเดียวกัน
UDF ที่ลงทะเบียนด้วย spark.udf.register() ในสมุดบันทึกเดียวกันได้รับการสนับสนุนสําหรับ CREATE ในทุกบริบท สําหรับการรีเฟรช LINEAGE เฉพาะบริบท PySpark เท่านั้นที่ได้รับการสนับสนุน เนื่องจากข้อกําหนด UDF ทํางานเป็นส่วนหนึ่งของการดําเนินการสมุดบันทึกตามกําหนดการ
from pyspark.sql import SparkSession
from pyspark.sql.types import StringType
from pyspark.sql.functions import expr
spark = SparkSession.builder.getOrCreate()
# UDF: Extract domain from email
def extract_email_domain(email):
if email is None or '@' not in email:
return None
return email.split('@')[1]
# Registration
spark.udf.register(
"udf_email_domain",
extract_email_domain,
StringType()
)
@fmlv.materialized_lake_view(
name="udf_testing_silver.mlv_high_value_customers",
comment="High-value customers identified by UDF criteria",
table_properties={"delta.enableChangeDataFeed": "true"}
)
def mlv_high_value_customers():
return spark.sql("""
SELECT
c.customer_id,
c.name,
c.email,
udf_email_domain(c.email) as email_domain,
c.segment,
c.lifetime_value,
total_transactions.total_amount,
total_transactions.txn_count
FROM udf_testing_bronze.customers c
INNER JOIN (
SELECT
customer_id,
SUM(amount) as total_amount,
COUNT(*) as txn_count
FROM udf_testing_bronze.transactions
WHERE udf_is_positive(amount)
GROUP BY customer_id
HAVING SUM(amount) > 1000
) total_transactions ON c.customer_id = total_transactions.customer_id
WHERE udf_validate_customer(c.email, c.age)
AND c.segment IN ('premium', 'vip')
""")
print("✓ Created mlv_high_value_customers")
ห้องสมุดของบุคคลที่สาม - Pandas UDFs
ไลบรารีของบุคคลที่สาม เช่น Pandas UDF อนุญาตให้ใช้กฎคุณภาพของข้อมูลด้วยการประมวลผลแบบเวกเตอร์ เปิดใช้งานการตรวจสอบขั้นสูง เช่น ตรรกะทางธุรกิจแบบกําหนดเอง การตรวจสอบทางสถิติ หรือการตรวจหารูปแบบที่ไม่สามารถทําได้ด้วยฟังก์ชันในตัว สิ่งนี้ช่วยสร้างข้อจํากัดด้านคุณภาพข้อมูลที่ปรับขนาดได้และนํากลับมาใช้ใหม่ได้ในระหว่างการสร้างและรีเฟรช MLV
import fmlv
from pyspark.sql.types import StructType, StructField, IntegerType, TimestampType, StringType, DoubleType, BooleanType
from datetime import datetime
import pandas as pd
from pyspark.sql.types import BooleanType
def pandas_check_impl(val):
# Reject if value < median of [100, 200, 300]
return val >= pd.Series([100, 200, 300]).median()
spark.udf.register("pandas_check", pandas_check_impl, BooleanType())
@fmlv.materialized_lake_view(
name="silver.pyspark_from_two_sqlmlv_inner_pandas",
comment="PySpark MLV INNER JOIN using pandas-based constraint and DROP violations"
)
@fmlv.check(
name="dq_pandas_check",
condition="pandas_check(l3)",
action="DROP"
)
def pyspark_from_two_sqlmlv_inner_pandas():
# Define the function
# Register the function as a Spark UDF
# Read source tables
df1 = spark.table("silver.base_sqlmlv")
df2 = spark.table("silver.base_sqlmlv")
# Rename columns for unique join
df_left = df1.select([df1[col].alias(f"l{i+1}") for i, col in enumerate(df1.columns)])
df_right = df2.select([df2[col].alias(f"r{i+1}") for i, col in enumerate(df2.columns)])
# Perform INNER JOIN
df = df_left.join(df_right, df_left.l1 == df_right.r1, "inner")
return df
df = spark.table("silver.pyspark_from_two_sqlmlv_inner_pandas")
# All amounts should be >= median (200)
assert all(df.select("l3").rdd.map(lambda r: r[0] >= 200).collect()), "Unexpected low-value rows found"
print("PySpark MLV INNER JOIN with pandas-based DQ DROP passed")
ไลบรารีที่กําหนดเอง – Python Wheel (.whl)
ฟังก์ชันที่บรรจุเป็นไฟล์ jar หรือไฟล์ล้อสามารถติดตั้งบนคลัสเตอร์ Fabric (ผ่านการตั้งค่าสภาพแวดล้อม) และใช้ในข้อกําหนด MLV CREATE และ LINEAGE REFRESH ได้รับการสนับสนุนสําหรับบริบท PySpark ดู จัดการไลบรารีแบบกําหนดเองในสภาพแวดล้อม Fabric สําหรับรายละเอียดเพิ่มเติม
%%pyspark
import fmlv
from pyspark.sql.types import StructType, StructField, IntegerType, TimestampType, StringType, DoubleType, BooleanType
from datetime import datetime
from pyspark.sql.types import BooleanType
from custom_dq_lib import threshold_check
def custom_check_impl(val):
return threshold_check(val, threshold=200)
spark.udf.register("custom_check", custom_check_impl, BooleanType())
@fmlv.materialized_lake_view(
name="silver.pyspark_from_two_sqlmlv_inner_custom_whl",
comment="PySpark MLV INNER JOIN using custom DQ library and DROP violations",
replace=True
)
@fmlv.check(
name="dq_custom_check",
condition="custom_check(l3)",
action="DROP"
)
def pyspark_from_two_sqlmlv_inner_custom():
# Wrap the custom function as Spark UDF
# Read source tables
df1 = spark.table("silver.base_sqlmlv")
df2 = spark.table("silver.base_sqlmlv")
# Rename columns for unique join
df_left = df1.select([df1[col].alias(f"l{i+1}") for i, col in enumerate(df1.columns)])
df_right = df2.select([df2[col].alias(f"r{i+1}") for i, col in enumerate(df2.columns)])
# Perform INNER JOIN
df = df_left.join(df_right, df_left.l1 == df_right.r1, "inner")
return df
df = spark.table("silver.pyspark_from_two_sqlmlv_inner_custom_whl")
# All amounts should be >= threshold (200)
assert all(df.select("l3").rdd.map(lambda r: r[0] >= 200).collect()), "Unexpected low-value rows found"
print("PySpark MLV INNER JOIN with custom DQ library passed")
ฟังก์ชันข้อมูลผู้ใช้ Fabric
ฟังก์ชันข้อมูลผู้ใช้ Fabric (UDF) ถูกกําหนดและจัดการจากส่วนกลางในพื้นที่ทํางาน Fabric พร้อมใช้งานสําหรับโน้ตบุ๊กหรือไปป์ไลน์ใดๆ โดยไม่จําเป็นต้องลงทะเบียนใหม่ต่อเซสชัน ทําให้เหมาะสําหรับไปป์ไลน์ MLV ที่ใช้งานจริง คุณลักษณะนี้ได้รับการสนับสนุนเฉพาะในบริบท PySpark สําหรับทั้ง CREATE และการรีเฟรช LINEAGE เรียนรู้เพิ่มเติมเกี่ยวกับฟังก์ชันข้อมูลผู้ใช้ที่นี่ สําหรับข้อมูลเพิ่มเติม โปรดดู ภาพรวมฟังก์ชันข้อมูลผู้ใช้ Fabric
%%pyspark
import fmlv
from pyspark.sql import functions as F
from notebookutils import udf
myFuncs = udf.getFunctions("UserDataFunction_1")
def add_greeting_column(df):
pdf = df.toPandas()
pdf["greeting"] = pdf["name"].apply(lambda n: myFuncs.hello_fabric(n))
return spark.createDataFrame(pdf)
import fmlv
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, DoubleType
# -------------------------------------------------------------
# Define base PySpark MLV (no SQL)
# -------------------------------------------------------------
@fmlv.materialized_lake_view(
name="silver.base_pysparkmlv",
comment="Base MLV created using PySpark",
replace=True
)
def base_pysparkmlv():
schema = StructType([
StructField("id", IntegerType(), True),
StructField("name", StringType(), True),
StructField("amount", DoubleType(), True),
StructField("country", StringType(), True)
])
data = [
(1, "Alice", 100.0, "US"),
(2, "Bob", 200.0, "UK"),
(3, "Charlie", 300.0, "UK")
]
return spark.createDataFrame(data, schema)
@fmlv.materialized_lake_view(
name="silver.mlv_udfn_null_test",
replace=False
)
@fmlv.check(
name="null_check",
condition="greeting IS NOT NULL",
action="FAIL"
)
def mlv_udfn_null_test():
df = spark.table("silver.base_pysparkmlv")
return add_greeting_column(df)