如何通过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
Levelfrom 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
tenacityor 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
breakstatement after detecting failures.
内容的提问来源于stack exchange,提问作者Justine
相关产品推荐
相关产品推荐

