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

如何将DynamoDB表导入PySpark?附JSON Line实现示例

直接从DynamoDB表读取数据到PySpark的方案

嗨,我来帮你实现直接以DynamoDB表作为数据源的PySpark操作,不用先导出到S3再读啦!下面是一步步的具体做法:

1. 先确保环境依赖到位

首先,你的PySpark环境需要能访问AWS DynamoDB的相关依赖包。如果是通过spark-submit启动任务,记得加上对应的依赖包参数(版本可以根据你的Spark和AWS SDK版本调整):

spark-submit --packages org.apache.hadoop:hadoop-aws:3.3.4,com.amazonaws:aws-java-sdk-dynamodb:1.12.500 your_script.py

如果是在EMR、Databricks这类托管环境,通常已经预装了这些依赖,不用额外添加。

2. 配置AWS访问权限

接下来要让Spark能访问你的DynamoDB表。如果是在AWS环境(比如EC2、EMR)运行,推荐用IAM角色赋予权限,不用硬编码凭证;如果是本地测试,可以在代码里配置(注意不要把凭证提交到代码仓库!):

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("ReadDynamoDB") \
    .getOrCreate()

# 配置AWS凭证(本地测试用,生产环境优先用IAM角色)
spark.conf.set("spark.hadoop.fs.s3a.access.key", "YOUR_AWS_ACCESS_KEY")
spark.conf.set("spark.hadoop.fs.s3a.secret.key", "YOUR_AWS_SECRET_KEY")
spark.conf.set("spark.hadoop.fs.s3a.region", "YOUR_AWS_REGION")  # 比如us-east-1

3. 定义和原需求一致的Schema

和你之前从S3读JSON时的Schema完全一样,直接复用就行:

from pyspark.sql.types import StructType, StructField, TimestampType, StringType, ArrayType, DoubleType

schema_table = StructType([
    StructField('created_at', TimestampType(), True),
    StructField('id', StringType(), True),
    StructField('matches', ArrayType(
        StructType([
            StructField('code', StringType(), True),
            StructField('prefix', StringType(), True),
            StructField('price', DoubleType(), True),
        ]), True
    ), True),
])

4. 直接读取DynamoDB表

现在就可以直接指定数据源为DynamoDB,传入表名和定义好的Schema了:

# 读取DynamoDB表
table = spark.read \
    .format("dynamodb") \
    .option("dynamodb.table.name", "YOUR_DYNAMODB_TABLE_NAME") \
    .schema(schema_table) \
    .load()

# 验证一下数据
table.show()
table.printSchema()

一些额外的注意事项

  • 数据类型映射:DynamoDB的String类型可以直接对应Spark的StringType,数值类型对应DoubleType/IntegerType。如果你的created_at在DynamoDB里是ISO格式的字符串,Spark的TimestampType会自动解析;如果是Unix时间戳,可能需要用from_unixtime函数转换。
  • 大表读取优化:如果你的DynamoDB表数据量很大,可以设置dynamodb.input.size参数调整每次读取的条目数(比如option("dynamodb.input.size", "1000")),或者开启并行扫描option("dynamodb.scan.parallelism", "4")来提高读取效率。
  • 权限问题:确保执行Spark的角色/用户有DynamoDB的Scan权限(因为Spark读取DynamoDB默认用Scan操作)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:34:00