【问题标题】:How to read data from S3 bucket using pyspark on local machine?如何在本地机器上使用 pyspark 从 S3 存储桶中读取数据?
【发布时间】:2021-09-22 08:08:38
【问题描述】:

我正在尝试使用 pyspark 从本地机器上的 S3 存储桶中读取数据。我从某个网站借了代码。当我提交代码时,它显示以下错误:

Traceback (most recent call last):
  File "C:\Users\Admin\Desktop\sparkcode.py", line 77, in <module>
    s3_df=spark.read.csv("s3a://bucket_name/dummy.csv",header=True,inferSchema=True)
  File "C:\spark\spark-3.1.2-bin-hadoop3.2\python\lib\pyspark.zip\pyspark\sql\readwriter.py", line 737, in csv
  File "C:\spark\spark-3.1.2-bin-hadoop3.2\python\lib\py4j-0.10.9-src.zip\py4j\java_gateway.py", line 1304, in __call__
  File "C:\spark\spark-3.1.2-bin-hadoop3.2\python\lib\pyspark.zip\pyspark\sql\utils.py", line 111, in deco
  File "C:\spark\spark-3.1.2-bin-hadoop3.2\python\lib\py4j-0.10.9-src.zip\py4j\protocol.py", line 326, in get_return_value
py4j.protocol.Py4JJavaError: An error occurred while calling o33.csv.
: java.lang.IllegalArgumentException
        at java.base/java.util.concurrent.ThreadPoolExecutor.<init>(ThreadPoolExecutor.java:1293)
        at java.base/java.util.concurrent.ThreadPoolExecutor.<init>(ThreadPoolExecutor.java:1215)
        at org.apache.hadoop.fs.s3a.S3AFileSystem.initialize(S3AFileSystem.java:280)
        at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:3303)
        at org.apache.hadoop.fs.FileSystem.access$200(FileSystem.java:124)
        at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:3352)
        at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:3320)
        at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:479)
        at org.apache.hadoop.fs.Path.getFileSystem(Path.java:361)
        at org.apache.spark.sql.execution.streaming.FileStreamSink$.hasMetadata(FileStreamSink.scala:46)
        at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:377)
        at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:325)
        at org.apache.spark.sql.DataFrameReader.$anonfun$load$3(DataFrameReader.scala:307)
        at scala.Option.getOrElse(Option.scala:189)
        at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:307)
        at org.apache.spark.sql.DataFrameReader.csv(DataFrameReader.scala:795)
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
        at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.base/java.lang.reflect.Method.invoke(Method.java:566)
        at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
        at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
        at py4j.Gateway.invoke(Gateway.java:282)
        at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
        at py4j.commands.CallCommand.execute(CallCommand.java:79)
        at py4j.GatewayConnection.run(GatewayConnection.java:238)
        at java.base/java.lang.Thread.run(Thread.java:834)

我正在使用以下命令执行:spark-submit sparkcode.py

sparkcode.py的内容是

import pyspark
from pyspark.sql import SparkSession
from pyspark import SparkContext, SparkConf

import os
os.environ['PYSPARK_SUBMIT_ARGS'] = '-- packages com.amazonaws:aws-java-sdk:1.7.4,org.apache.hadoop:hadoop-aws:2.7.3 pyspark-shell'

#spark configuration
conf = SparkConf().set('spark.executor.extraJavaOptions','-Dcom.amazonaws.services.s3.enableV4=true').set('spark.driver.extraJavaOptions','-Dcom.amazonaws.services.s3.enableV4=true').set("spark.hadoop.fs.s3a.multipart.size", 104857600).setAppName('pyspark_aws').setMaster('local[*]')
sc=SparkContext(conf=conf)
sc.setSystemProperty('com.amazonaws.services.s3.enableV4', 'true')

print("modules imported")

hadoopConf = sc._jsc.hadoopConfiguration()
hadoopConf.set('spark.hadoop.fs.s3a.access.key', 'access_key')
hadoopConf.set('spark.hadoop.fs.s3a.secret.key', 'secret_key')
hadoopConf.set('spark.hadoop.fs.s3a.endpoint', 's3-us-east-2.amazonaws.com')
hadoopConf.set('spark.hadoop.fs.s3a.impl', 'org.apache.hadoop.fs.s3a.S3AFileSystem')

spark=SparkSession(sc)

s3_df=spark.read.csv("s3a://bucket_name/dummy.csv",header=True,inferSchema=True)
print(s3_df.show())

我怎样才能摆脱这个错误?

【问题讨论】:

  • 你能提供所有的日志吗?从这里调试应用程序真的很难
  • @RobertKossendey 我已经编辑了问题并添加了完整的日志。

标签: amazon-web-services apache-spark hadoop amazon-s3 pyspark


【解决方案1】:

您能否尝试以下示例 pyspark (spark_s3_integration.py) 代码:

from __future__ import print_function
import sys

from pyspark.conf import SparkConf
from pyspark.sql import SparkSession
from pyspark.sql import Row

if __name__ == "__main__":
    if len(sys.argv) != 4:
        print("Usage  : spark_s3_integration.py <AWS_ACCESS_KEY_ID> <AWS_SECRET_ACCESS_KEY> <BUCKET_NAME>", file=sys.stderr)
        print("Example: spark_s3_integration.py ranga_aws_access_key ranga_aws_secret_key ranga-spark-s3-bkt", file=sys.stderr)
        exit(-1)

    awsAccessKey = sys.argv[1]
    awsSecretKey = sys.argv[2]
    bucketName = sys.argv[3]

    conf = (
        SparkConf()
            .setAppName("PySpark S3 Integration Example")
            .set("spark.hadoop.fs.s3a.access.key", awsAccessKey)
            .set("spark.hadoop.fs.s3a.secret.key", awsSecretKey)
            .set("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")
            .set("spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version", "2")
            .set("spark.speculation", "false")
            .set("spark.hadoop.mapreduce.fileoutputcommitter.cleanup-failures.ignored", "true")
            .set("fs.s3a.experimental.input.fadvise", "random")
            .setIfMissing("spark.master", "local")
    )

    # Creating the SparkSession object
    spark = SparkSession.builder.config(conf=conf).getOrCreate()
    print("SparkSession Created successfully")

    Employee = Row("id", "name", "age", "salary")
    employee1 = Employee(1, "Ranga", 32, 245000.30)
    employee2 = Employee(2, "Nishanth", 2, 345000.10)
    employee3 = Employee(3, "Raja", 32, 245000.86)
    employee4 = Employee(4, "Mani", 14, 45000.00)

    employeeData = [employee1, employee2, employee3, employee4]
    employeeDF = spark.createDataFrame(employeeData)
    employeeDF.printSchema()
    employeeDF.show()

    # Define the s3 destination path
    s3_dest_path = "s3a://" + bucketName + "/employees"
    print("s3 destination path "+s3_dest_path)

    # Write the data as Orc
    employeeOrcPath = s3_dest_path + "/employee_orc"
    employeeDF.write.mode("overwrite").format("orc").save(employeeOrcPath)

    # Read the employee orc data
    employeeOrcData = spark.read.format("orc").load(employeeOrcPath);
    employeeOrcData.printSchema()
    employeeOrcData.show()

    # Write the data as Parquet
    employeeParquetPath = s3_dest_path + "/employee_parquet"
    employeeOrcData.write.mode("overwrite").format("parquet").save(employeeParquetPath)

    spark.stop()
    print("SparkSession stopped")

源代码:

https://github.com/rangareddy/ranga_spark_experiments/blob/master/spark_s3_integration/spark_s3_integration.py

【讨论】:

  • 1.删除“fs.s3a.impl”声明;这是一个堆栈溢出火花迷信。如果您不相信我,请尝试 2​​. 使用 S3A 提交者来确保安全和性能
  • @stevel 我确认从本地在 S3 上写入文件不需要配置,但我想知道 S3A 提交者是什么。 Spark 仍然很新,它的许多概念对我来说还很陌生。
猜你喜欢
  • 2022-01-11
  • 2022-01-13
  • 2021-10-25
  • 1970-01-01
  • 1970-01-01
  • 2021-03-13
  • 1970-01-01
  • 2018-11-07
  • 1970-01-01
相关资源
最近更新 更多