15TB Snappy压缩Parquet数据从S3迁移至DynamoDB的1-5天实现方案咨询
S3到DynamoDB大规模数据迁移方案(15TB Parquet)
方案选择:优先使用Spark
Spark的并行处理能力和DynamoDB官方集成的连接器,比Hive更适合大规模数据的快速迁移,尤其是需要灵活添加字段(如「插入日期」)和处理写入重试的场景。
集群配置建议
实例类型与节点数量(按时间要求分档)
1天内完成迁移
- 主节点:
m5.xlarge(4vCPU,16GB内存),仅负责集群管理,无需高配置 - 核心节点:12台
r5.2xlarge(8vCPU,32GB内存,内存优化型适合Parquet文件的内存读取) - 实例模式:核心节点用按需实例保证稳定性,可选添加2-4台Spot任务节点扩容
3-5天内完成迁移
- 主节点:
m5.xlarge(同上) - 核心节点:6台
m5.2xlarge(8vCPU,32GB内存,性价比更高) - 实例模式:核心节点按需,任务节点可选Spot
Spark配置参数(提交任务时设置)
spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 6g \ --executor-cores 2 \ --num-executors 48 \ # 12台核心节点×4个executor/节点 --driver-memory 8g \ --conf spark.dynamodb.throughput.write.percent=0.8 \ --conf spark.dynamodb.retry.max=10 \ --conf spark.dynamodb.retry.delay.base=50 \ your-migration-script.py
- 每个核心节点分配4个executor(32GB内存留4GB给系统,每个executor用6GB+2vCPU,匹配8vCPU的节点规格)
- 调整
num-executors:核心节点数×4即可
DynamoDB写入优化(含指数退避重试)
核心配置参数
通过Spark-DynamoDB连接器的内置参数实现指数退避和流量控制:
spark.dynamodb.throughput.write.percent=0.8:限制写入吞吐量占按需模式突发容量的80%,避免触发限流spark.dynamodb.retry.max=10:最大重试次数,覆盖大部分临时限流场景spark.dynamodb.retry.delay.base=50:初始重试延迟50ms,后续按指数递增(50ms→100ms→200ms…)spark.dynamodb.batch.write.size=25:按DynamoDB单批最大限制批量写入,减少API调用次数
热点规避
由于原S3分区键与DynamoDB分区键/排序键不匹配,写入前需确保DynamoDB的分区键值分布均匀:
- 若目标表分区键本身是均匀分布的(如UUID、散列值),直接写入即可
- 若分区键存在倾斜,可在Spark中对数据做shuffle打散后再写入
迁移步骤示例(Python版Spark代码)
from pyspark.sql import SparkSession from pyspark.sql.functions import current_date # 初始化Spark会话 spark = SparkSession.builder \ .appName("S3ParquetToDynamoDB") \ .getOrCreate() # 读取S3上的Snappy压缩Parquet数据 df = spark.read.parquet("s3://your-source-bucket/path/to/parquet-data/") # 添加「插入日期」列(取当前日期) df_with_insert_date = df.withColumn("插入日期", current_date()) # 写入DynamoDB df_with_insert_date.write \ .format("dynamodb") \ .option("dynamodb.table.name", "your-target-dynamodb-table") \ .option("dynamodb.throughput.write.percent", "0.8") \ .option("dynamodb.retry.max", "10") \ .option("dynamodb.retry.delay.base", "50") \ .option("dynamodb.batch.write.size", "25") \ .mode("append") \ .save() # 停止Spark会话 spark.stop()
额外注意事项
- 使用EMR最新稳定版本(如emr-6.10.0),自带兼容的Spark和DynamoDB连接器,无需手动安装依赖
- 开启EMR日志存储到S3,方便排查写入失败或性能瓶颈问题
- 若使用Spot实例,配置任务重试策略,确保节点中断后任务能自动恢复
内容的提问来源于stack exchange,提问作者dba
相关产品推荐
相关产品推荐

