通过Glue PySpark作业在Redshift执行CTAS语句并实现并行执行
问题描述
- 我在AWS Glue中集成Redshift与Apache Spark,希望直接向Redshift发送
CREATE TABLE AS SELECT * FROM source_table这类语句执行,无需将数据读入DataFrame。但当前代码似乎是在Glue Catalog中创建表,而非Redshift,求简便实现方法。 - 需要运行大量此类建表语句,如何用PySpark并行化请求并同时执行?
附上原测试代码:
from pyspark.sql import SparkSession # Replace with your Redshift JDBC URL, user, password, and driver class redshift_jdbc_url = "jdbc:redshift://your-redshift-endpoint:5439/your-database?user=your-user&password=your-password" redshift_user = "your-user" redshift_password = "your-password" redshift_driver = "com.amazon.redshift.jdbc.Driver" # Initialize Spark session spark = SparkSession.builder \ .appName("RedshiftCTASExample") \ .config("spark.jars", "path/to/RedshiftJDBC42-2.0.0.4.jar") \ .getOrCreate() # Define the SQL query query = """ CREATE TABLE new_table AS ( SELECT * FROM event ) """ # Execute the query in Redshift spark.read \ .format("jdbc") \ .option("url", redshift_jdbc_url) \ .option("user", redshift_user) \ .option("password", redshift_password) \ .option("driver", redshift_driver) \ .option("query", query) \ .load() # Stop the Spark session spark.stop()
解决方案
一、直接在Redshift执行CTAS语句
你当前代码用spark.read.jdbc()执行CTAS的问题在于:这个API是用来读取数据返回DataFrame的,当执行无返回结果的DDL/DML时,Spark会尝试解析结果并触发Glue Catalog元数据同步,导致误创建Glue表。正确做法是直接通过JDBC连接执行无返回的SQL语句,以下是两种常用方式:
方法1:使用JDBC原生连接(Glue环境推荐)
通过PySpark调用Java的JDBC连接,直接执行语句,不涉及Spark DataFrame读取:
from pyspark.sql import SparkSession # Glue环境中可直接使用glueContext.spark_session,无需手动初始化 spark = SparkSession.builder \ .appName("RedshiftDirectCTAS") \ .config("spark.jars", "path/to/RedshiftJDBC42-2.0.0.4.jar") \ .getOrCreate() # 提取Redshift连接参数 jdbc_url = "jdbc:redshift://your-redshift-endpoint:5439/your-database" user = "your-user" password = "your-password" driver = "com.amazon.redshift.jdbc.Driver" # 获取Java JDBC连接并执行CTAS jvm = spark._jvm driver_class = jvm.Class.forName(driver) connection = driver_class.getConnect(jdbc_url, user, password) statement = connection.createStatement() try: query = "CREATE TABLE new_table AS SELECT * FROM event" statement.execute(query) finally: statement.close() connection.close() spark.stop()
注意:Glue环境中JDBC Jar包可通过作业的「依赖JAR」配置上传,无需指定本地路径。
方法2:使用Spark 3.3+的executeQuery API
如果你的Spark版本≥3.3.0,可使用专门用于无返回SQL的executeQuery方法:
spark.read.format("jdbc") \ .option("url", jdbc_url) \ .option("user", user) \ .option("password", password) \ .option("driver", driver) \ .executeQuery("CREATE TABLE new_table AS SELECT * FROM event")
二、并行执行大量建表语句
结合Redshift的并发查询能力(受集群WLM队列配置限制),可通过以下两种方式实现并行:
方法1:Spark并行化RDD执行
将所有CTAS语句放入列表,通过Spark的parallelize创建RDD,每个分区任务独立建立JDBC连接执行:
from pyspark.sql import SparkSession def execute_ctas(query): # 每个任务必须独立初始化JDBC连接,不能共享连接对象 jvm = spark._jvm driver = "com.amazon.redshift.jdbc.Driver" jvm.Class.forName(driver) conn = jvm.DriverManager.getConnection(jdbc_url, user, password) stmt = conn.createStatement() try: stmt.execute(query) finally: stmt.close() conn.close() return f"Success: {query}" # 批量CTAS语句列表 ctas_queries = [ "CREATE TABLE table1 AS SELECT * FROM source1", "CREATE TABLE table2 AS SELECT * FROM source2", "CREATE TABLE table3 AS SELECT * FROM source3" # 更多语句... ] # 根据Redshift WLM并发数设置并行度(例如8) spark.sparkContext.parallelize(ctas_queries, numSlices=8).map(execute_ctas).collect()
方法2:Python线程池执行(适合小批量)
如果语句数量不多,可使用Python线程池实现并行:
from concurrent.futures import ThreadPoolExecutor def execute_ctas_thread(query): jvm = spark._jvm driver = "com.amazon.redshift.jdbc.Driver" jvm.Class.forName(driver) conn = jvm.DriverManager.getConnection(jdbc_url, user, password) stmt = conn.createStatement() try: stmt.execute(query) finally: stmt.close() conn.close() return f"Success: {query}" # 线程数匹配Redshift并发能力 with ThreadPoolExecutor(max_workers=8) as executor: results = list(executor.map(execute_ctas_thread, ctas_queries)) print(results)
关键注意事项
- Redshift WLM配置:并行前需确认Redshift的WLM队列允许足够并发查询,避免排队或报错。
- 连接隔离:每个并行任务/线程必须独立创建JDBC连接,JDBC连接不是线程安全的。
- 权限控制:确保Glue作业角色拥有Redshift的
redshift:ExecuteQuery权限,以及JDBC Jar包的访问权限。
内容的提问来源于stack exchange,提问作者L Xandor
相关产品推荐
相关产品推荐

