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

从AWS Glue Catalog向Redshift加载数据的技术需求问询

问题解决方案

1. 实现加载前后的自定义操作(如截断维度表)

可以通过Glue的JDBC连接执行自定义SQL来实现前置/后置操作,比如截断维度表。示例代码如下:

# 定义前置截断SQL
truncate_sql = f"TRUNCATE TABLE {schema}.{dimension_table};"

# 获取Glue连接中的JDBC信息
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
jdbc_conf = glueContext.extract_jdbc_conf("redshift_connection")
jdbc_url = jdbc_conf["url"]
conn_props = {
    "user": jdbc_conf["user"],
    "password": jdbc_conf["password"],
    "driver": "com.amazon.redshift.jdbc42.Driver"
}

# 执行前置SQL
conn = spark._jvm.java.sql.DriverManager.getConnection(jdbc_url, conn_props["user"], conn_props["password"])
conn.createStatement().execute(truncate_sql)
conn.close()

# 原有数据加载逻辑
dyf = glueContext.create_dynamic_frame.from_catalog(
    database=catalog_db,
    table_name=f"{table}",
    push_down_predicate=f"day = {day}"
)
glueContext.write_dynamic_frame.from_jdbc_conf(
    frame=dyf,
    catalog_connection="redshift_connection",
    connection_options={
        "database": database,
        "dbtable": f"{schema}.{table}"
    },
    redshift_tmp_dir= "s3://your-tmp-bucket/path/"
)

# 可选:执行后置操作(比如写入审计日志)
post_sql = "INSERT INTO audit.load_log VALUES ('dimension_load', CURRENT_TIMESTAMP);"
conn = spark._jvm.java.sql.DriverManager.getConnection(jdbc_url, conn_props["user"], conn_props["password"])
conn.createStatement().execute(post_sql)
conn.close()

可以把执行SQL的逻辑封装成函数,针对不同表判断是否需要执行前置/后置操作,提升扩展性。

2. JDBC URL直连 vs Glue Catalog连接的对比

支持JDBC URL搭配用户名密码的方式

直接用JDBC URL的代码示例:

dyf = glueContext.create_dynamic_frame.from_catalog(
    database=catalog_db,
    table_name=f"{table}",
    push_down_predicate=f"day = {day}"
)

glueContext.write_dynamic_frame.from_jdbc_conf(
    frame=dyf,
    catalog_connection=None,
    connection_options={
        "url": "jdbc:redshift://your-redshift-endpoint:5439/your-db",
        "user": "your-username",
        "password": "your-password",
        "dbtable": f"{schema}.{table}"
    },
    redshift_tmp_dir= "s3://your-tmp-bucket/path/"
)

两种方式优劣对比

  • Glue Catalog连接(推荐)
    • 优势:敏感信息(密码)存在Glue连接中,无需硬编码;支持IAM角色认证,不用手动管理用户名密码;连接信息统一维护,修改时无需改代码,符合AWS安全最佳实践。
    • 劣势:依赖Glue连接资源,调试时需确保配置正确。
  • JDBC URL直连
    • 优势:配置直接,不依赖Glue连接,适合快速测试。
    • 劣势:用户名密码易泄露(硬编码或散落在配置中);连接信息分散,不利于统一维护;无法利用IAM认证,需手动管理凭证。

3. 提前创建Schema的可选方案

如果不需要Glue自动建表,有两种实现方式:

方式一:Redshift提前手动建表

先在Redshift中执行CREATE TABLE语句定义好表结构,然后在Glue写数据时指定写入模式,避免自动建表:

glueContext.write_dynamic_frame.from_jdbc_conf(
    frame=dyf,
    catalog_connection="redshift_connection",
    connection_options={
        "database": database,
        "dbtable": f"{schema}.{table}",
        "mode": "append",  # 或"overwrite",根据需求选择
        "preactions": ""  # 清空Glue默认的自动建表前置动作
    },
    redshift_tmp_dir= "s3://your-tmp-bucket/path/"
)

方式二:Glue代码中执行建表SQL

在数据加载前,先执行建表SQL确保结构存在,再加载数据:

# 定义建表SQL(可根据Glue Catalog的表结构生成)
create_table_sql = f"""
CREATE TABLE IF NOT EXISTS {schema}.{table} (
    id INT,
    name VARCHAR(100),
    day DATE,
    amount DECIMAL(10,2)
);
"""

# 执行建表SQL(复用之前的JDBC连接逻辑)
jdbc_conf = glueContext.extract_jdbc_conf("redshift_connection")
jdbc_url = jdbc_conf["url"]
conn_props = {
    "user": jdbc_conf["user"],
    "password": jdbc_conf["password"],
    "driver": "com.amazon.redshift.jdbc42.Driver"
}

spark = SparkSession.builder.getOrCreate()
conn = spark._jvm.java.sql.DriverManager.getConnection(jdbc_url, conn_props["user"], conn_props["password"])
conn.createStatement().execute(create_table_sql)
conn.close()

# 执行数据加载,此时不会自动建表
glueContext.write_dynamic_frame.from_jdbc_conf(
    frame=dyf,
    catalog_connection="redshift_connection",
    connection_options={
        "database": database,
        "dbtable": f"{schema}.{table}",
        "mode": "append"
    },
    redshift_tmp_dir= "s3://your-tmp-bucket/path/"
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 05:22:45