You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 12:57:53