1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
| # 完整的ETL管道
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.spark.operators.spark_submit import SparkSubmitOperator
from datetime import datetime, timedelta
def build_data_pipeline():
"""构建数据管道"""
default_args = {
'owner': 'data-team',
'depends_on_past': False,
'start_date': datetime(2024, 1, 1),
'email_on_failure': True,
'email_on_retry': False,
'retries': 2,
'retry_delay': timedelta(minutes=5)
}
dag = DAG(
'data_lake_pipeline',
default_args=default_args,
description='数据湖ETL管道',
schedule_interval='0 2 * * *', # 每天凌晨2点
catchup=False,
tags=['datalake', 'etl']
)
# 任务1:从Kafka摄取数据到Bronze层
kafka_to_bronze = SparkSubmitOperator(
task_id='kafka_to_bronze',
application='jobs/ingest_kafka_to_bronze.py',
name='ingest-kafka-to-bronze',
conn_id='spark_default',
conf={
'spark.dynamicAllocation.enabled': 'true',
'spark.dynamicAllocation.maxExecutors': '10',
'spark.executor.memory': '4g',
'spark.executor.cores': '2'
},
dag=dag
)
# 任务2:清洗数据到Silver层
bronze_to_silver = SparkSubmitOperator(
task_id='bronze_to_silver',
application='jobs/bronze_to_silver.py',
name='bronze-to-silver',
conf={
'spark.sql.adaptive.enabled': 'true',
'spark.sql.adaptive.coalescePartitions.enabled': 'true'
},
dag=dag
)
# 任务3:构建聚合表到Gold层
silver_to_gold = SparkSubmitOperator(
task_id='silver_to_gold',
application='jobs/silver_to_gold.py',
name='silver-to-gold',
dag=dag
)
# 任务4:数据质量检查
data_quality_check = PythonOperator(
task_id='data_quality_check',
python_callable=run_quality_checks,
dag=dag
)
# 任务5:数据血缘更新
update_lineage = PythonOperator(
task_id='update_lineage',
python_callable=update_data_lineage,
dag=dag
)
# 任务6:发送报告
send_report = PythonOperator(
task_id='send_report',
python_callable=send_pipeline_report,
dag=dag
)
# 设置任务依赖
kafka_to_bronze >> bronze_to_silver >> silver_to_gold
silver_to_gold >> data_quality_check >> update_lineage >> send_report
return dag
# 主函数
def run_quality_checks(**context):
"""运行数据质量检查"""
from quality_checker import DataQualityChecker
checker = DataQualityChecker()
# 检查空值率
null_checks = checker.check_nullity("gold.fact_orders")
# 检查数据范围
range_checks = checker.check_ranges("gold.fact_orders", {
"total_amount": (0, 1000000),
"order_count": (0, 10000)
})
# 检查数据完整性
integrity_checks = checker.check_referential_integrity(
"gold.fact_orders",
{"dim_products": ["product_id"], "dim_users": ["user_id"]}
)
# 生成质量报告
report = {
"null_checks": null_checks,
"range_checks": range_checks,
"integrity_checks": integrity_checks,
"overall_score": calculate_overall_score([
null_checks, range_checks, integrity_checks
])
}
if report["overall_score"] < 0.9:
raise Exception(f"Data quality score too low: {report['overall_score']}")
return report
def calculate_overall_score(checks):
"""计算总体质量分数"""
scores = []
for check in checks:
if isinstance(check, dict):
scores.append(check.get("score", 1.0))
return sum(scores) / len(scores) if scores else 1.0
|