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

通过Glue PySpark作业在Redshift执行CTAS语句并实现并行执行

问题描述
  1. 我在AWS Glue中集成Redshift与Apache Spark,希望直接向Redshift发送CREATE TABLE AS SELECT * FROM source_table这类语句执行,无需将数据读入DataFrame。但当前代码似乎是在Glue Catalog中创建表,而非Redshift,求简便实现方法。
  2. 需要运行大量此类建表语句,如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 04:43:26