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

Databricks Unity Catalog中组和视图创建代码的优化咨询

Databricks Unity Catalog 组与视图创建代码优化

需求说明

我有一段可正常运行的代码,用于在Databricks工作区中先检查组和视图是否已存在,再创建对应的组与视图。现希望对该代码进行优化,需注意我使用的是Databricks Unity Catalog,并非所有PySpark SQL属性都受支持。

原代码

target_database = "main.deltalake_db"
distinctprojects = spark.sql("SELECT DISTINCT requestProjectUrlName FROM main.deltalake_db.Customers_Log_Table_test")

for project in distinctprojects.collect():
    project_name = project['requestProjectUrlName']
    view_name = f"view_{project_name.replace(' ', '_').replace('-', '_')}"
    groupname = f"users_{view_name}"
    
    # Check if group already exists
    group_exists = False
    for group in spark.sql("SHOW GROUPS").collect():
        if 'name' in group:
            group_name = group['name']
        else:
            group_name = group['groupName']
        if group_name == groupname:
            group_exists = True
            print(f"Group {groupname} already exists.")
            break
    
    if not group_exists:
        create_group = f"""
        CREATE GROUP {groupname} with USER `xyz@test.com`
        """
        try: 
            spark.sql(create_group)            
        except Exception as e:
            print(f"Error creating group {groupname}: {e}")
        else:
            print(f"Group {groupname} created successfully.")
    
    # Check if view already exists
    view_exists = False
    for table in spark.catalog.listTables(target_database):
        if table.name == view_name:
            view_exists = True
            print(f"View {view_name} already exists.")
            break
    
    if not view_exists:
        create_view_query = f"""
        CREATE OR REPLACE VIEW {target_database}.{view_name} AS
        SELECT * FROM {target_database}.Customers_Log_Table
        WHERE 
          CASE
            WHEN is_member('{groupname}') AND requestProjectUrlName = '{project_name}' THEN TRUE 
            ELSE FALSE
      END;
    """ 
    try: 
        spark.sql(create_view_query)
        
    except Exception as e:
        print(f"Error creating view {view_name}: {e}")
    else:
        print(f"View {view_name} created successfully.")

优化后的代码

target_database = "main.deltalake_db"
admin_user = "`xyz@test.com`"

# 1. 批量获取所有需要处理的项目名称
distinct_projects_df = spark.sql("SELECT DISTINCT requestProjectUrlName FROM main.deltalake_db.Customers_Log_Table_test")
project_names = [row['requestProjectUrlName'] for row in distinct_projects_df.collect()]

# 2. 预处理生成组名和视图名,确保符合Unity Catalog命名规范
def sanitize_name(name):
    # 替换空格、连字符、点为下划线,仅保留合法字符
    return name.replace(' ', '_').replace('-', '_').replace('.', '_')

project_mappings = []
for proj_name in project_names:
    sanitized_proj = sanitize_name(proj_name)
    view_name = f"view_{sanitized_proj}"
    group_name = f"users_{view_name}"
    project_mappings.append({
        "original_proj": proj_name,
        "view_name": view_name,
        "group_name": group_name
    })

# 3. 批量查询已存在的组,减少多次查询开销
if project_mappings:
    group_names_to_check = [mapping['group_name'] for mapping in project_mappings]
    existing_groups_df = spark.sql(f"""
        SELECT name FROM system.information_schema.groups 
        WHERE name IN ({', '.join([f"'{g}'" for g in group_names_to_check])})
    """)
    existing_groups = set(row['name'] for row in existing_groups_df.collect())

    # 4. 处理组创建
    for mapping in project_mappings:
        group_name = mapping['group_name']
        if group_name not in existing_groups:
            create_group_sql = f"CREATE GROUP {group_name} WITH USER {admin_user}"
            try:
                spark.sql(create_group_sql)
                print(f"Group {group_name} created successfully.")
            except Exception as e:
                print(f"Error creating group {group_name}: {str(e)}")
        else:
            print(f"Group {group_name} already exists.")

    # 5. 批量查询已存在的视图
    view_names_to_check = [mapping['view_name'] for mapping in project_mappings]
    existing_views_df = spark.sql(f"""
        SELECT table_name FROM system.information_schema.views 
        WHERE table_catalog = 'main' AND table_schema = 'deltalake_db' 
          AND table_name IN ({', '.join([f"'{v}'" for v in view_names_to_check])})
    """)
    existing_views = set(row['table_name'] for row in existing_views_df.collect())

    # 6. 处理视图创建
    for mapping in project_mappings:
        view_name = mapping['view_name']
        group_name = mapping['group_name']
        original_proj = mapping['original_proj']
        
        if view_name not in existing_views:
            # 简化WHERE条件,避免冗余CASE逻辑
            create_view_sql = f"""
            CREATE VIEW {target_database}.{view_name} AS
            SELECT * FROM {target_database}.Customers_Log_Table
            WHERE requestProjectUrlName = '{original_proj}' AND is_member('{group_name}')
            """
            try:
                spark.sql(create_view_sql)
                print(f"View {view_name} created successfully.")
            except Exception as e:
                print(f"Error creating view {view_name}: {str(e)}")
        else:
            print(f"View {view_name} already exists.")

优化说明

  • 批量查询降本增效:原代码每次循环全量查询所有组/视图,优化后一次性批量查询目标对象,大幅减少Spark作业执行次数与资源消耗。
  • 统一命名规范:新增sanitize_name函数,确保生成的组名、视图名符合Unity Catalog命名规则(仅允许字母、数字、下划线),避免非法字符导致创建失败。
  • 系统表替代兼容查询:使用system.information_schema.groups和system.information_schema.views系统表替代SHOW GROUPS与listTables,结果结构更稳定,无需处理字段名兼容问题。
  • 简化逻辑提升可读性:将视图的CASE条件简化为直接的AND逻辑,代码更清晰;拆分预处理、组处理、视图处理模块,可维护性更强。
  • 降低SQL注入风险:通过命名 sanitize 处理,减少非法字符注入的可能性。

内容的提问来源于stack exchange,提问作者ZZZSharePoint

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 13:19:59