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

如何通过Snowpark按级别顺序执行DDL语句表?

Modified Snowpark Python Code for Level-Ordered DDL Execution

Here's the adjusted code that executes DDL statements from highest Level to lowest, runs same-level statements in parallel, and includes basic error handling:

import snowflake.snowpark as snowpark
from snowflake.snowpark.exceptions import SnowparkSQLException

def execute_ddl_by_level(session: snowpark.Session, table_name: str):
    # Load DDL table into DataFrame
    ddl_table = session.table(table_name)
    df = ddl_table.to_pandas()
    
    # Sort records by Level descending (highest first)
    df_sorted = df.sort_values(by='Level', ascending=False)
    
    # Group DDLs by their Level
    grouped_ddl = df_sorted.groupby('Level')
    
    # Process each Level from highest to lowest
    for level, group in sorted(grouped_ddl, key=lambda x: x[0], reverse=True):
        print(f"Processing Level {level} DDLs...")
        async_jobs = []
        
        # Submit all DDLs in the current Level as async jobs
        for _, row in group.iterrows():
            ddl_cmd = row['DDL']
            try:
                job = session.sql(ddl_cmd).collect_nowait()
                async_jobs.append((ddl_cmd, job))
                print(f"Submitted: {ddl_cmd}")
            except Exception as e:
                print(f"Failed to submit DDL: {ddl_cmd}. Error: {str(e)}")
        
        # Wait for all jobs in the Level to finish and check results
        level_has_failures = False
        for ddl_cmd, job in async_jobs:
            try:
                # Wait for job completion; raises exception if execution fails
                job.result()
                print(f"Success: {ddl_cmd}")
            except SnowparkSQLException as e:
                print(f"FAILED: {ddl_cmd}. Error: {str(e)}")
                level_has_failures = True
            except Exception as e:
                print(f"Unexpected error: {ddl_cmd}. Error: {str(e)}")
                level_has_failures = True
        
        # Stop execution if current Level has failures (lower levels depend on these objects)
        if level_has_failures:
            print(f"Level {level} has failed DDLs. Aborting further execution.")
            break
    
    print("Execution process finished.")

# Run the function
tableName = 'DATABASE.SCHEMA.DDL_TABLE'
execute_ddl_by_level(new_session, tableName)

Key Features

  • Level Ordering: Sorts DDLs by Level from highest to lowest, ensuring dependent objects (like Level 0 views) are created only after their dependencies (Level 5-1 tables).
  • Parallel Execution: Same-level DDLs are submitted as asynchronous jobs, running in parallel to optimize execution time.
  • Error Handling: Waits for all jobs in a level to complete, checks for failures, and stops execution immediately if any DDL in a higher level fails—preventing wasted effort on dependent lower-level objects.
  • Transparent Logging: Logs every submission, success, and failure to help track progress and debug issues.

Optional Enhancements

  • Retry Logic: Add retries for transient errors (e.g., warehouse contention) using a library like tenacity or custom loop logic.
  • Concurrency Throttling: If you have hundreds of DDLs in one level, limit the number of concurrent jobs to avoid overloading your Snowflake warehouse.
  • Partial Execution: If you need to proceed even with failed DDLs (not recommended for dependent objects), remove the break statement after detecting failures.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 09:28:20