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

PySpark集成Amazon DynamoDB无数据返回问题排查及替代方案咨询

PySpark连接DynamoDB无数据问题排查与替代方案

问题背景

使用com.audienceproject:spark-dynamodb_2.12:1.1.2库连接DynamoDB后,仅返回空表结构,无实际数据。

现有代码

import os
from pyspark.sql import SparkSession

os.environ["AWS_ACCESS_KEY"] ="myAccessKey" 
os.environ["AWS_SECRET_ACCESS_KEY"] = "mySecretKey"
spark = SparkSession.builder \
    .appName("SQL Server to PySpark") \
    .config("spark.jars.packages","com.audienceproject:spark-dynamodb_2.12:1.1.2") \
    .getOrCreate()

df  = spark.read.option("tableName", "test_db") \
                     .format("dynamodb") \
                     .load()
df.show()

关键现象

  • 依赖包加载正常,无报错
  • 输出为空表:
++
| |
++
++

排查步骤

  • 权限验证:确认AWS密钥拥有目标表的dynamodb:Scan权限,用AWS CLI执行aws dynamodb scan --table-name test_db --region your-region测试能否获取数据。
  • 区域配置:该库默认使用us-east-1,若表在其他区域,需添加区域参数:
    df = spark.read.option("tableName", "test_db") \
                    .option("region", "your-region") \
                    .format("dynamodb") \
                    .load()
    
  • 版本兼容:Spark 3.5.0与1.1.2版本可能存在兼容性问题,升级到最新版com.audienceproject:spark-dynamodb_2.12:1.3.0重试。
  • 表数据确认:通过AWS控制台检查test_db表是否确实存在数据。

替代连接方案

方案1:AWS Glue官方连接器

兼容性更强,官方维护:

spark = SparkSession.builder \
    .appName("DynamoDB-PySpark") \
    .config("spark.jars.packages", "com.amazonaws:aws-glue-dynamodb:1.0.0") \
    .getOrCreate()

df = spark.read.format("dynamodb") \
    .option("dynamodb.table.name", "test_db") \
    .option("dynamodb.region", "your-region") \
    .option("aws.access.key.id", "myAccessKey") \
    .option("aws.secret.access.key", "mySecretKey") \
    .load()
df.show()

方案2:boto3读取转DataFrame

适合小数据量场景:

import boto3
from pyspark.sql import Row

dynamodb = boto3.resource('dynamodb', 
                          region_name='your-region',
                          aws_access_key_id='myAccessKey',
                          aws_secret_access_key='mySecretKey')
table = dynamodb.Table('test_db')
response = table.scan()
items = response['Items']

# 转换为Spark DataFrame
df = spark.createDataFrame([Row(**item) for item in items])
df.show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 14:08:10