嵌入式实战项目教学:从环境监控终端到系统能力沉淀
2026/10/2 16:43:11
frompyspark.sqlimportSparkSessionfrompyspark.sql.functionsimportcol,coalesce,trim,when,lit,sumfrompyspark.sql.typesimportStringType,NumericType# 在 Databricks 中,spark 会话通常已经存在,无需重新创建# 如果需要显式创建,使用:# spark = SparkSession.builder.getOrCreate()# 配置参数database_name="your_database"# 替换为实际数据库名result_list=[]# 获取数据库中的所有表和视图(包括 Delta 表)tables=spark.catalog.listTables(database_name)fortableintables:table_name=table.name full_table_name=f"{database_name}.{table_name}"try:# 读取表或视图(Delta 表和普通表均可)df=spark.table(full_table_name)# 快速判断是否为空表,避免不必要的缓存total_count=df.count()iftotal_count==0:continue# 缓存数据以便多次引用(实际上我们只需扫描一次,缓存并非必需)df.cache()# 构建所有字段的聚合表达式(一次性计算)agg_exprs=[]field_meta=[]# 用于记录每个字段的类型和名称forfieldindf.schema.fields:col_name=field.name col_type=field.dataTypeifisinstance(col_type,StringType):# 字符串类型:统计 null 或 trim 后为空字符串modified_col=trim(coalesce(col(col_name),lit("")))condition=(modified_col==lit(""))count_expr=sum(when(condition,1).otherwise(0)).alias(f"cnt_{col_name}")elifisinstance(col_type,NumericType):# 数值类型:统计 null 或零值modified_col=coalesce(col(col_name),lit(0))condition=(modified_col==lit(0))count_expr=sum(when(condition,1).otherwise(0)).alias(f"cnt_{col_name}")else:# 其他类型:仅统计 nullcondition=col(col_name).isNull()count_expr=sum(when(condition,1).otherwise(0)).alias(f"cnt_{col_name}")agg_exprs.append(count_expr)field_meta.append((col_name,str(col_type)))# 执行一次聚合,获取所有字段的统计值stats_row=df.agg(*agg_exprs).collect()[0]# 整理结果forcol_name,col_typeinfield_meta:stat_count=stats_row[f"cnt_{col_name}"]percentage=round((stat_count/total_count)*100,2)iftotal_count>0else0.0result_list.append((database_name,table_name,col_name,col_type,stat_count,total_count,float(percentage)))df.unpersist()# 释放缓存exceptExceptionase:print(f"Error processing table{table_name}:{str(e)}")continue# 创建结果 DataFrameresult_columns=["database_name","table_name","column_name","column_type","stat_count","total_rows","percentage"]result_df=spark.createDataFrame(result_list,result_columns)# 显示结果result_df.show(truncate=False)# 可选:将结果保存到 Delta 表# result_df.write.format("delta").mode("overwrite").saveAsTable("your_audit_table")| 原代码 | 修改后 | 原因 |
|---|---|---|
显式创建SparkSession并启用 Hive | 直接使用 Databricks 内置spark | Databricks 已预配置,无需额外初始化 |
对每个字段分别执行df.agg() | 构建所有字段聚合表达式,一次df.agg() | 减少表扫描次数,大幅提高性能 |
缓存df后多次扫描 | 缓存后仅一次聚合扫描 | 配合优化,降低缓存开销 |
spark.catalog.listTables() | 同样使用,能返回表和视图 | 无需修改,原生支持 Delta 表和视图 |
| 异常处理与结果收集 | 保持不变 | 逻辑通用 |
database_name,如果使用默认数据库,可设为"default"。spark.table()可以读取持久化视图,但临时视图不会出现在listTables中。如需处理临时视图,可额外指定名称列表。df.count()和聚合仍会触发一次完整扫描,这是必要的。如果表极大,可考虑分批处理或使用抽样估算。