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

如何使用AWS Glue一次性实现跨账号多表AWS DynamoDB数据迁移

问题原因

你遇到的序列化报错核心原因是:GlueContext、SparkContext属于Driver端专属对象,仅允许在Driver进程中调用。你尝试将这类对象广播到Worker节点、或是在分布式算子(flatMap/foreach)内调用相关方法,会触发序列化失败。分布式算子的逻辑会被分发到Worker节点运行,Worker节点无法访问Driver侧的Spark上下文对象。

正确实现方案

不需要使用分布式算子处理多表逻辑,所有表的读写调度都在Driver侧执行即可,有两种常用实现方式:

方案1:串行循环实现(适合表数量较少的场景)

直接循环迁移表列表,逻辑简单稳定,代码如下:

import boto3
import sys
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from awsglue.utils import getResolvedOptions

args = getResolvedOptions(sys.argv, ["JOB_NAME"])
sc = SparkContext.getOrCreate()
glue_context = GlueContext(sc)
job = Job(glue_context)
job.init(args["JOB_NAME"], args)

# 配置参数,可根据实际情况修改
region = "your-region"
ddb_split = 10
# 定义要迁移的表映射:键为源表名,值为目标表名
table_mapping = {
    "source_table1": "target_table1",
    "source_table2": "target_table2"
}
# 跨账号凭证,建议通过STS AssumeRole获取,不要硬编码
ddb_key = "your-access-key"
ddb_secret = "your-secret-key"
ddb_token = "your-session-token"

# 循环处理每个表
for source_table, target_table in table_mapping.items():
    print(f"开始迁移表:{source_table} -> {target_table}")
    # 读源表
    dyf = glue_context.create_dynamic_frame_from_options(
        connection_type="dynamodb",
        connection_options={
            "dynamodb.region": region,
            "dynamodb.splits": str(ddb_split),
            "dynamodb.throughput.read.percent": "0.5", # 建议调低避免打满源表RCU
            "dynamodb.input.tableName": source_table
        }
    )
    # 写目标表
    glue_context.write_dynamic_frame_from_options(
        frame=dyf,
        connection_type="dynamodb",
        connection_options={
            "dynamodb.region": region,
            "dynamodb.output.tableName": target_table,
            "dynamodb.awsAccessKeyId": ddb_key,
            "dynamodb.awsSecretAccessKey": ddb_secret,
            "dynamodb.awsSessionToken": ddb_token
        }
    )
    print(f"表{source_table}迁移完成")

print("所有表迁移完成!")
job.commit()

方案2:多线程并行实现(适合表数量多、需要提升迁移效率的场景)

DynamoDB读写属于IO密集型操作,用Python多线程在Driver端并行执行多个表的迁移任务,既可以提升效率,也不会触发序列化问题:

import boto3
import sys
from concurrent.futures import ThreadPoolExecutor, as_completed
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from awsglue.utils import getResolvedOptions

args = getResolvedOptions(sys.argv, ["JOB_NAME"])
sc = SparkContext.getOrCreate()
glue_context = GlueContext(sc)
job = Job(glue_context)
job.init(args["JOB_NAME"], args)

# 配置参数
region = "your-region"
ddb_split = 10
table_mapping = {
    "source_table1": "target_table1",
    "source_table2": "target_table2",
    "source_table3": "target_table3"
}
ddb_key = "your-access-key"
ddb_secret = "your-secret-key"
ddb_token = "your-session-token"
# 并行度,根据源表RCU、目标表WCU配置调整,避免超限
max_workers = 2

def migrate_table(source_table, target_table):
    print(f"开始迁移表:{source_table} -> {target_table}")
    dyf = glue_context.create_dynamic_frame_from_options(
        connection_type="dynamodb",
        connection_options={
            "dynamodb.region": region,
            "dynamodb.splits": str(ddb_split),
            "dynamodb.throughput.read.percent": "0.5",
            "dynamodb.input.tableName": source_table
        }
    )
    glue_context.write_dynamic_frame_from_options(
        frame=dyf,
        connection_type="dynamodb",
        connection_options={
            "dynamodb.region": region,
            "dynamodb.output.tableName": target_table,
            "dynamodb.awsAccessKeyId": ddb_key,
            "dynamodb.awsSecretAccessKey": ddb_secret,
            "dynamodb.awsSessionToken": ddb_token
        }
    )
    return f"表{source_table}迁移完成"

# 提交并行任务
with ThreadPoolExecutor(max_workers=max_workers) as executor:
    futures = [executor.submit(migrate_table, src, dst) for src, dst in table_mapping.items()]
    for future in as_completed(futures):
        print(future.result())

print("所有表迁移完成!")
job.commit()
注意事项
  • 提前在目标AWS账号创建对应DynamoDB表,保证分区键、排序键、索引结构和源表一致
  • 跨账号凭证推荐通过boto3 sts assume_role接口临时获取,不要硬编码AKSK到代码中
  • 调整dynamodb.throughput.read.percent参数和并行度,避免超过源表读容量、目标表写容量限制

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 22:48:03