【问题标题】:Combining delta io and excel reading结合delta io和excel阅读
【发布时间】:2021-11-22 03:10:18
【问题描述】:

当使用没有 delta 的 com.crealytics:spark-excel_2.12:0.14.0 时:

spark = SparkSession.builder.appName("Word Count")
.config("spark.jars.packages", "com.crealytics:spark-excel_2.12:0.14.0")
.getOrCreate()

df = spark.read.format("com.crealytics.spark.excel")
.option("header", "true")
.load(path2)

它工作正常,我可以很好地读取 excel 文件。但是使用 configure_spark_with_delta_pip 创建会话:

builder = SparkSession.builder.appName("transaction")
.config("spark.jars.packages", "com.crealytics:spark-excel_2.12:0.14.0")
.config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
.config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")

spark = configure_spark_with_delta_pip(builder).getOrCreate()

给我以下错误:

Py4JJavaError:调用 o139.load 时出错。 : java.lang.ClassNotFoundException:找不到数据源: com.crealytics.spark.excel。请在以下位置找到包裹 http://spark.apache.org/third-party-projects.html 在 org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:692) 在 org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSourceV2(DataSource.scala:746) 在 org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:265) 在 org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:239) 在 java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native 方法)在 java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) 在 java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 在 java.base/java.lang.reflect.Method.invoke(Method.java:566) 在 py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) 在 py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) 在 py4j.Gateway.invoke(Gateway.java:282) 在 py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) 在 py4j.commands.CallCommand.execute(CallCommand.java:79) 在 py4j.GatewayConnection.run(GatewayConnection.java:238) 在 java.base/java.lang.Thread.run(Thread.java:829) 原因: java.lang.ClassNotFoundException: com.crealytics.spark.excel.DefaultSource 在 java.base/java.net.URLClassLoader.findClass(URLClassLoader.java:471) 在 java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:589) 在 java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:522) 在 org.apache.spark.sql.execution.datasources.DataSource$.$anonfun$lookupDataSource$5(DataSource.scala:666) 在 scala.util.Try$.apply(Try.scala:213) 在 org.apache.spark.sql.execution.datasources.DataSource$.$anonfun$lookupDataSource$4(DataSource.scala:666) 在 scala.util.Failure.orElse(Try.scala:224) 在 org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:666) ... 14 更多

为什么?我该如何避免这种情况?

【问题讨论】:

    标签: apache-spark pyspark delta-lake


    【解决方案1】:

    您收到此错误是因为 configure_spark_with_delta_pip 覆盖/替换了您的配置属性 spark.jars.packages 要导入的适当 delta Lake 包。因此,您的包 com.crealytics:spark-excel_2.12:0.14.0 可能不可用/导入。查看来自source code here的sn-p

        scala_version = "2.12"
        maven_artifact = f"io.delta:delta-core_{scala_version}:{delta_version}"
    
        return spark_session_builder.config("spark.jars.packages", maven_artifact) 
    

    不幸的是,此时Builder 不允许我们检索现有的配置属性或SparkConf 对象以在调用getOrCreate 创建或触发会话之前动态调整这些属性。

    方法 1

    要解决此问题,您可以自己检索适当的 delta 包,类似于 configure_spark_with_delta_pip 的做法,例如。

    
    import importlib_metadata
    delta_version = importlib_metadata.version("delta_spark")
    scala_version = "2.12"
    delta_package = f"io.delta:delta-core_{scala_version}:{delta_version}"
    
    builder = SparkSession.builder.appName("transaction")
    .config("spark.jars.packages", f"com.crealytics:spark-excel_2.12:0.14.0,{delta_package}")
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
    
    spark = configure_spark_with_delta_pip(builder).getOrCreate()
    
    
    

    方法2

    要解决此问题,您可以在使用 configure_spark_with_delta_pip 应用 delta 包后创建 spark 会话。在此之后,您可以使用更新的配置属性触发重新启动 spark 会话。

    例如。

    builder = SparkSession.builder.appName("transaction")
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
    
    spark = configure_spark_with_delta_pip(builder).getOrCreate()
    
    builder = SparkSession.builder.appName("transaction")
    .config("spark.jars.packages", "com.crealytics:spark-excel_2.12:0.14.0")
    
    spark = builder.getOrCreate()
    

    由于两个 spark 会话具有相同的 appNamegetOrCreate 将检索现有的 spark 会话,但也会应用新配置。此行为记录在 here

    如果返回现有的 SparkSession,配置选项 此构建器中指定的将应用于现有的 SparkSession。

    >>> s2 = SparkSession.builder.config("k2", "v2").getOrCreate()
    >>> s1.conf.get("k1") == s2.conf.get("k1") 
    True
    >>> s1.conf.get("k2") == s2.conf.get("k2") 
    True
    

    让我知道这是否适合你。

    【讨论】:

      猜你喜欢
      • 2017-09-01
      • 1970-01-01
      • 2021-01-04
      • 2018-06-24
      • 1970-01-01
      • 2011-02-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多