从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
相关产品推荐
相关产品推荐

