如何使用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
相关产品推荐
相关产品推荐

