如何通过AWS Glue作业向Redshift批量上传多张表
如何用AWS Glue作业批量将多张表加载到Redshift
当然可以批量处理多张表!既然你不熟悉Python,我会给你一个简单易改的示例脚本——基于你原来的自动生成脚本修改,只需要做少量调整就能适配多表场景。
核心思路
把单表的处理逻辑封装成一个可复用的函数,然后定义你要处理的表名列表,循环调用这个函数即可。这样不用重复写相同的代码,也方便后续添加更多表。
批量处理示例脚本
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job ## @params: [TempDir, JOB_NAME] args = getResolvedOptions(sys.argv, ['TempDir','JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) # -------------------------- 自定义配置部分 -------------------------- # 定义你要处理的表名列表(从Glue Data Catalog里的sampledb数据库) TABLES_TO_PROCESS = ["abs", "table2", "table3"] # 替换成你的实际表名 GLUE_DATABASE = "sampledb" REDSHIFT_CONNECTION_NAME = "redshift" REDSHIFT_DATABASE = "dbmla" # 如果不同表的字段映射不一样,可以用字典存储每个表的映射 # 示例:TABLE_MAPPINGS = { # "abs": [("value", "int", "value", "int"), ("sex", "string", "sex", "string"), ...], # "table2": [("field1", "string", "field1", "string"), ...] # } # 这里先使用你原来的abs表映射作为示例,如果其他表映射不同,记得修改 DEFAULT_MAPPING = [ ("value", "int", "value", "int"), ("sex", "string", "sex", "string"), ("age", "string", "age", "string"), ("highest year of school completed", "string", "highest year of school completed", "string"), ("state", "string", "state", "string"), ("region type", "string", "region type", "string"), ("lga 2011", "string", "lga 2011", "string"), ("frequency", "string", "frequency", "string"), ("time", "string", "time", "string") ] # ------------------------------------------------------------------- # 定义处理单张表的函数 def process_single_table(table_name): print(f"开始处理表: {table_name}") # 1. 从Glue Catalog读取数据 datasource = glueContext.create_dynamic_frame.from_catalog( database=GLUE_DATABASE, table_name=table_name, transformation_ctx=f"datasource_{table_name}" ) # 2. 应用字段映射(如果不同表有不同映射,这里可以改成TABLE_MAPPINGS[table_name]) applymapping = ApplyMapping.apply( frame=datasource, mappings=DEFAULT_MAPPING, transformation_ctx=f"applymapping_{table_name}" ) # 3. 处理字段类型选择 resolvechoice = ResolveChoice.apply( frame=applymapping, choice="make_cols", transformation_ctx=f"resolvechoice_{table_name}" ) # 4. 去除空值字段 dropnullfields = DropNullFields.apply( frame=resolvechoice, transformation_ctx=f"dropnullfields_{table_name}" ) # 5. 写入Redshift datasink = glueContext.write_dynamic_frame.from_jdbc_conf( frame=dropnullfields, catalog_connection=REDSHIFT_CONNECTION_NAME, connection_options={ "dbtable": table_name, # Redshift里的目标表名,和Glue表名一致,可按需修改 "database": REDSHIFT_DATABASE }, redshift_tmp_dir=args["TempDir"], transformation_ctx=f"datasink_{table_name}" ) print(f"表 {table_name} 处理完成") # 循环处理所有表 for table in TABLES_TO_PROCESS: process_single_table(table) # 提交作业 job.commit()
如何调整使用
- 修改表名列表:把
TABLES_TO_PROCESS里的内容替换成你实际需要处理的Glue表名,比如["abs", "customer", "orders"]。 - 字段映射调整:如果不同表的字段映射规则不一样,把
DEFAULT_MAPPING改成字典形式的TABLE_MAPPINGS,然后在process_single_table函数里用mappings=TABLE_MAPPINGS[table_name]即可。 - Redshift目标表名:如果Redshift里的表名和Glue表名不一样,修改
connection_options里的"dbtable"值,比如改成f"public.{table_name}"指定schema。
注意事项
- 确保Glue Data Catalog里的所有目标表都已经存在(你已经完成了爬取步骤,这一步应该没问题)。
- 确认你的Glue作业角色有访问S3临时目录、Glue Catalog和Redshift的权限。
- 如果Redshift里还没有创建对应的目标表,可以在
connection_options里添加"create_table_options": "FILLFACTOR 100"来让Glue自动创建表(前提是字段映射正确)。
内容的提问来源于stack exchange,提问作者beni
相关产品推荐
相关产品推荐

