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

如何通过AWS Glue自动获取Document DB所有集合并导出数据?

问题描述

我需要实现以下目标:

  • 连接Document DB集群
  • 连接集群中的特定数据库
  • 获取该数据库下所有集合(表)的列表
  • 从所有表中提取500行数据
  • 将每个表导出为.csv文件并存储到S3存储桶

目前的痛点是:创建爬虫时需要硬编码大量信息(比如集合名称),无法为数据库中的所有集合逐一执行操作,求自动化流程。

已完成的尝试

通过AWS Glue处理单个集合的流程:

  1. 创建自定义连接器
  2. 通过该连接创建爬虫,获取表的schema
  3. 从数据目录创建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()

配置要点

  1. 作业参数:在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/)
  2. 依赖配置:

    • 脚本依赖pymongo,需在Glue作业中添加该库:可通过创建Lambda层上传pymongo包,或在Glue作业的Python库路径指定包含该库的S3路径。
  3. 权限配置:

    • Glue作业角色需具备DocumentDB的连接权限(安全组开放27017端口给Glue)
    • Glue作业角色需具备S3输出桶的读写权限

方法二:基于Glue爬虫的批量处理

如果偏好使用Glue数据目录,可通过以下流程自动化:

  1. 创建批量爬虫:配置Glue爬虫连接DocumentDB时,选择"包含所有集合",让爬虫自动发现数据库下的所有表并注册到数据目录。
  2. 编写遍历脚本:在Glue作业中遍历数据目录中的所有目标表,读取每个表的500行数据并导出到S3。

这种方式无需手动连接DocumentDB,但依赖Glue爬虫的调度和表发现能力,适合需要长期维护表Schema的场景。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 17:20:01