如何使用AWS SessionToken在PySpark中读取S3数据?
使用AWS SessionToken在PySpark中读取S3数据
我来帮你搞定这个问题~看了你给出的代码,目前还没配置AWS临时会话凭证,这才是连接S3的关键。下面是完整的解决方案和需要注意的细节:
1. 核心配置:添加SessionToken到SparkConf
你需要在SparkConf里设置三个AWS临时凭证参数,分别是fs.s3a.access.key、fs.s3a.secret.key和fs.s3a.session.token,对应你的AWS临时会话里的三个凭证值。
2. 修正后的完整代码
import os # 注意:hadoop-aws版本要和你的Spark依赖的Hadoop版本匹配!比如Spark2.x配2.7.x,Spark3.x配3.x os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages "org.apache.hadoop:hadoop-aws:2.7.3" pyspark-shell' from pyspark import SparkConf, SparkContext conf = SparkConf() \ .setMaster("local[2]") \ .setAppName("pyspark-unittests") \ .set("spark.sql.parquet.compression.codec", "snappy") \ # 这里填入你的临时凭证 .set("fs.s3a.access.key", "你的临时AccessKey") \ .set("fs.s3a.secret.key", "你的临时SecretKey") \ .set("fs.s3a.session.token", "你的SessionToken") sc = SparkContext(conf=conf) # 读取S3上的文件 s3File = sc.textFile("s3a://myrepo/test.csv") print(s3File.count()) # 记得关闭上下文 sc.stop()
3. 必看的注意事项
- 版本一定要匹配:
hadoop-aws的版本必须和Spark内置的Hadoop版本对应,不然会出现各种奇怪的兼容性错误。比如Spark2.4.x默认用Hadoop2.7.x,所以你选2.7.3没问题;如果是Spark3.1+,建议用hadoop-aws:3.2.0。 - 用对协议:一定要用
s3a://而不是s3://或者s3n://,s3a是Hadoop对S3的最新实现,支持更多功能,性能也更好。 - 别硬编码凭证:代码里直接写凭证不安全,建议通过环境变量读取,比如:
这样运行前只要把三个环境变量设置好就行。import os conf.set("fs.s3a.access.key", os.getenv('AWS_ACCESS_KEY_ID')) conf.set("fs.s3a.secret.key", os.getenv('AWS_SECRET_ACCESS_KEY')) conf.set("fs.s3a.session.token", os.getenv('AWS_SESSION_TOKEN')) - 权限要到位:确保你的临时会话凭证拥有读取目标S3桶和文件的权限(至少要有
s3:GetObject权限),不然会报权限错误。
4. 另一种方式:用Hadoop配置文件
如果不想在代码里写配置,也可以修改Spark的core-site.xml配置文件,添加以下内容:
<property> <name>fs.s3a.access.key</name> <value>你的临时AccessKey</value> </property> <property> <name>fs.s3a.secret.key</name> <value>你的临时SecretKey</value> </property> <property> <name>fs.s3a.session.token</name> <value>你的SessionToken</value> </property>
然后启动Spark的时候指定这个配置文件的路径就可以了。
内容的提问来源于stack exchange,提问作者Jared
相关产品推荐
相关产品推荐

