如何通过AWS Glue自动获取Document DB所有集合并导出数据?
问题描述
我需要实现以下目标:
- 连接Document DB集群
- 连接集群中的特定数据库
- 获取该数据库下所有集合(表)的列表
- 从所有表中提取500行数据
- 将每个表导出为.csv文件并存储到S3存储桶
目前的痛点是:创建爬虫时需要硬编码大量信息(比如集合名称),无法为数据库中的所有集合逐一执行操作,求自动化流程。
已完成的尝试
通过AWS Glue处理单个集合的流程:
- 创建自定义连接器
- 通过该连接创建爬虫,获取表的schema
- 从数据目录创建ETL任务,运行脚本提取500行数据并导出为.csv格式上传至S3存储桶
对应的单集合处理脚本:
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 args = getResolvedOptions(sys.argv, ["JOB_NAME"]) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args["JOB_NAME"], args) read_mongo_options = { "uri": "mongodb://docdb-xyz.ap-south-1.docdb.amazonaws.com:27017", "database": "test", "collection": "collectionname", "username": "user", "ssl": "true", "password": "password123" } # Script generated for node Data Catalog table DataCatalogtable_node1 = glueContext.create_sample_dynamic_frame_from_options( connection_type="mongodb", connection_options=read_mongo_options, num=500, transformation_ctx="DataCatalogtable_node1" ) applymapping1 = ApplyMapping.apply( frame = DataCatalogtable_node1, mappings=[ ("account number", "long", "account number", "long"), ("cust_sensitivity", "string", "cust_sensitivity", "string"), ("pan number", "string", "pan number", "string"), ("cust_loyalty_points", "int", "cust_loyalty_points", "int"), ("graduation_yrof_ind", "int", "graduation_yrof_ind", "int"), ("gender", "string", "gender", "string"), ("aadhaar_card_number", "double", "aadhaar_card_number", "double"), ("no_of_2_wheelers_owned", "int", "no_of_2_wheelers_owned", "int"), ( "mi_is_betterthan_competition", "string", "mi_is_betterthan_competition", "string", ), ("degree_pref_mulcar_family", "int", "degree_pref_mulcar_family", "int"), ("cust_retention_overall", "int", "cust_retention_overall", "int"), ("complaints_cnt", "int", "complaints_cnt", "int"), ("no_of_dependents", "int", "no_of_dependents", "int"), ("family_member_cnt", "int", "family_member_cnt", "int"), ("cust_retention_sale", "int", "cust_retention_sale", "int"), ("complaints_open_cnt", "int", "complaints_open_cnt", "int"), ( "upsell_crosssell_opportunity", "string", "upsell_crosssell_opportunity", "string", ), ("cust_segment", "string", "cust_segment", "string"), ("purpose_of_vehicle", "string", "purpose_of_vehicle", "string"), ("yrof_1st_job_after_grad", "int", "yrof_1st_job_after_grad", "int"), ("tot_mul_cars_owned_family", "int", "tot_mul_cars_owned_family", "int"), ("marital-status", "string", "marital-status", "string"), ("cust_mhi_sale", "int", "cust_mhi_sale", "int"), ("no_of_earning_members", "int", "no_of_earning_members", "int"), ("become_parent_last_1st_yr", "string", "become_parent_last_1st_yr", "string"), ("no_of_children", "int", "no_of_children", "int"), ("zipcode", "int", "zipcode", "int"), ("cust_mood", "string", "cust_mood", "string"), ("cust_type", "string", "cust_type", "string"), ("average_household_savings", "int", "average_household_savings", "int"), ("tot_nonmul_cars_owned_family", "int", "tot_nonmul_cars_owned_family", "int"), ( "didyou_evenbuy_nonmi_for_veh", "string", "didyou_evenbuy_nonmi_for_veh", "string", ), ("occupation", "string", "occupation", "string"), ("date of birth", "string", "date of birth", "string"), ("dob_children1", "string", "dob_children1", "string"), ("_id", "struct", "_id", "string"), ("email-id", "string", "email-id", "string"), ("tot_cars_owned_family", "int", "tot_cars_owned_family", "int"), ("dob_children2", "string", "dob_children2", "string"), ("svoc_id", "string", "svoc_id", "string"), ("mobile_number", "long", "mobile_number", "long"), ("cust_tot_clv", "int", "cust_tot_clv", "int"), ("personal_annual_income", "int", "personal_annual_income", "int"), ("cust_name", "string", "cust_name", "string"), ("location", "string", "location", "string"), ], transformation_ctx = "applymapping1" ) dyf_repartion = applymapping1.coalesce(1) # Script generated for node S3 bucket S3bucket_node3 = glueContext.write_dynamic_frame.from_options( frame=dyf_repartion, connection_type="s3", format="csv", connection_options={"path": "s3://ouputbucket", "partitionKeys": []}, transformation_ctx="S3bucket_node3", ) job.commit()
自动化解决方案
方法一:动态遍历集合的Glue脚本
直接在Glue作业中连接DocumentDB获取所有集合,批量处理导出,无需依赖爬虫。
核心思路
- 用
pymongo连接DocumentDB,获取目标数据库的集合列表 - 遍历每个集合,动态生成读取配置
- 用
ResolveChoice和Relationalize自动处理Schema,替代硬编码的字段映射 - 将每个集合的500行数据导出到S3的独立目录
完整自动化脚本
import sys import pymongo from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job # 解析作业参数 args = getResolvedOptions(sys.argv, ["JOB_NAME", "DOC_DB_URI", "DB_NAME", "USERNAME", "PASSWORD", "S3_OUTPUT_PATH"]) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args["JOB_NAME"], args) # 连接DocumentDB获取所有集合名称 client = pymongo.MongoClient( args["DOC_DB_URI"], username=args["USERNAME"], password=args["PASSWORD"], ssl=True ) db = client[args["DB_NAME"]] collections = db.list_collection_names() # 批量处理每个集合 for coll in collections: # 动态配置集合读取参数 read_options = { "uri": args["DOC_DB_URI"], "database": args["DB_NAME"], "collection": coll, "username": args["USERNAME"], "ssl": "true", "password": args["PASSWORD"] } # 读取500行样本数据 dyf = glueContext.create_sample_dynamic_frame_from_options( connection_type="mongodb", connection_options=read_options, num=500, transformation_ctx=f"dyf_{coll}" ) # 自动处理Schema冲突和嵌套结构 dyf_resolved = ResolveChoice.apply( frame=dyf, choice="make_cols", transformation_ctx=f"resolved_{coll}" ) dyf_flat = Relationalize.apply( frame=dyf_resolved, staging_path=f"{args['S3_OUTPUT_PATH']}/temp", name=f"root_{coll}", transformation_ctx=f"flat_{coll}" ) # 合并为单个CSV文件 dyf_repartition = dyf_flat.coalesce(1) # 导出到S3对应集合目录 glueContext.write_dynamic_frame.from_options( frame=dyf_repartition, connection_type="s3", format="csv", connection_options={ "path": f"{args['S3_OUTPUT_PATH']}/{coll}", "partitionKeys": [], "header": "true" }, transformation_ctx=f"write_{coll}" ) job.commit()
配置要点
作业参数:在Glue作业中配置以下参数(避免硬编码敏感信息):
DOC_DB_URI:DocumentDB集群连接URI(如mongodb://docdb-xyz.ap-south-1.docdb.amazonaws.com:27017)DB_NAME:目标数据库名(如test)USERNAME:数据库用户名PASSWORD:数据库密码S3_OUTPUT_PATH:S3输出根路径(如s3://your-output-bucket/docdb-exports/)
依赖配置:
- 脚本依赖
pymongo,需在Glue作业中添加该库:可通过创建Lambda层上传pymongo包,或在Glue作业的Python库路径指定包含该库的S3路径。
- 脚本依赖
权限配置:
- Glue作业角色需具备DocumentDB的连接权限(安全组开放27017端口给Glue)
- Glue作业角色需具备S3输出桶的读写权限
方法二:基于Glue爬虫的批量处理
如果偏好使用Glue数据目录,可通过以下流程自动化:
- 创建批量爬虫:配置Glue爬虫连接DocumentDB时,选择"包含所有集合",让爬虫自动发现数据库下的所有表并注册到数据目录。
- 编写遍历脚本:在Glue作业中遍历数据目录中的所有目标表,读取每个表的500行数据并导出到S3。
这种方式无需手动连接DocumentDB,但依赖Glue爬虫的调度和表发现能力,适合需要长期维护表Schema的场景。
内容的提问来源于stack exchange,提问作者rishab ajain445
相关产品推荐
相关产品推荐

