Azure Synapse统计Lake/SQL数据库表行数并存储到DataFrame
实现Azure Synapse中Lake/SQL数据库全表行数统计并存储为DataFrame
核心逻辑
针对Lake数据库(Delta表)优先通过版本历史统计各版本行数,无历史数据则统计当日数据;SQL数据库直接统计当前行数,最终将所有结果合并到DataFrame。以下是具体实现方案:
1. 环境与依赖
确保在Synapse Spark池中运行代码,导入必要库:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, LongType, DateType import pyspark.sql.functions as F
2. 遍历数据库与表
编写工具函数获取指定类型的所有数据库和表:
def get_all_tables(spark, db_type="lake"): tables = [] # 获取所有Lake数据库及Delta表 if db_type == "lake": dbs = spark.sql("SHOW DATABASES").select("databaseName").rdd.flatMap(lambda x: x).collect() for db in dbs: tbls = spark.sql(f"SHOW TABLES IN {db}").select("tableName").rdd.flatMap(lambda x: x).collect() for tbl in tbls: # 验证是否为Delta表 try: spark.sql(f"DESCRIBE DETAIL {db}.{tbl}").select("format").collect()[0][0] tables.append((db, tbl)) except: continue # 获取所有SQL池数据库及表 elif db_type == "sql": sql_pool_name = "your-sql-pool-name" # 替换为你的SQL池名称 dbs = spark.sql(f"SHOW DATABASES IN {sql_pool_name}").select("databaseName").rdd.flatMap(lambda x: x).collect() for db in dbs: tbls = spark.sql(f"SHOW TABLES IN {sql_pool_name}.{db}").select("tableName").rdd.flatMap(lambda x: x).collect() for tbl in tbls: tables.append((db, tbl, sql_pool_name)) return tables
3. Lake数据库(Delta表)行数统计
处理Delta表的版本历史与降级逻辑:
def process_lake_tables(spark, lake_tables): schema = StructType([ StructField("database_name", StringType(), True), StructField("table_name", StringType(), True), StructField("row_count", LongType(), True), StructField("version", StringType(), True), StructField("stat_date", DateType(), True) ]) result_df = spark.createDataFrame([], schema) for db, tbl in lake_tables: full_tbl_name = f"{db}.{tbl}" try: # 获取Delta表所有版本历史 history_df = spark.sql(f"DESCRIBE HISTORY {full_tbl_name}").select("version") versions = history_df.select("version").rdd.flatMap(lambda x: x).collect() # 遍历版本统计行数 for version in versions: count_df = spark.sql(f""" SELECT '{db}' AS database_name, '{tbl}' AS table_name, COUNT(*) AS row_count, '{version}' AS version, DATE(current_timestamp()) AS stat_date FROM {full_tbl_name} VERSION AS OF {version} """) result_df = result_df.union(count_df) except Exception as e: # 无版本历史时,统计当前版本数据 print(f"{full_tbl_name}无法获取版本历史,统计当前数据: {str(e)}") count_df = spark.sql(f""" SELECT '{db}' AS database_name, '{tbl}' AS table_name, COUNT(*) AS row_count, 'current' AS version, DATE(current_timestamp()) AS stat_date FROM {full_tbl_name} """) result_df = result_df.union(count_df) return result_df
4. SQL数据库(Synapse SQL池)行数统计
处理SQL池表的当前行数统计:
def process_sql_tables(spark, sql_tables): schema = StructType([ StructField("database_name", StringType(), True), StructField("table_name", StringType(), True), StructField("row_count", LongType(), True), StructField("version", StringType(), True), StructField("stat_date", DateType(), True) ]) result_df = spark.createDataFrame([], schema) for db, tbl, sql_pool in sql_tables: full_tbl_name = f"{sql_pool}.{db}.{tbl}" try: count_df = spark.sql(f""" SELECT '{db}' AS database_name, '{tbl}' AS table_name, COUNT(*) AS row_count, 'current' AS version, DATE(current_timestamp()) AS stat_date FROM {full_tbl_name} """) result_df = result_df.union(count_df) except Exception as e: print(f"{full_tbl_name}统计失败: {str(e)}") continue return result_df
5. 合并结果并落地
将Lake和SQL表的统计结果合并为最终DataFrame,支持每日调度积累历史:
if __name__ == "__main__": spark = SparkSession.builder.appName("TableRowCount").getOrCreate() # 获取所有目标表 lake_tables = get_all_tables(spark, db_type="lake") sql_tables = get_all_tables(spark, db_type="sql") # 处理并合并结果 lake_result_df = process_lake_tables(spark, lake_tables) sql_result_df = process_sql_tables(spark, sql_tables) final_df = lake_result_df.union(sql_result_df) # 可选:将结果写入历史表,每日追加积累 # final_df.write.mode("append").format("delta").saveAsTable("your_history_db.row_count_history") # 展示结果 final_df.show()
注意事项
- 权限配置:确保Spark池拥有Lake数据库(ADLS/Gen2)的Storage Blob Data Contributor权限,以及Synapse SQL池的SQL Contributor权限。
- 性能优化:大表全量
COUNT(*)耗时较长,可通过DESCRIBE DETAIL {full_tbl_name}获取Delta表内置的rowCount元数据替代全量统计。 - 调度配置:将代码封装为Synapse管道的Spark作业,配置每日定时触发,自动积累历史统计数据。
内容的提问来源于stack exchange,提问作者coding
相关产品推荐
相关产品推荐

