【问题标题】:How to refer deltalake tables in jupyter notebook using pyspark如何使用 pyspark 在 jupyter notebook 中引用 deltalake 表
【发布时间】:2020-01-04 12:45:42
【问题描述】:

我正在尝试使用Pyspark 开始使用DeltaLakes

为了能够使用 deltalake,我在 Anaconda shell-prompt 上调用 pyspark 作为 —

pyspark — packages io.delta:delta-core_2.11:0.3.0

这是来自 deltalake 的参考资料 — https://docs.delta.io/latest/quick-start.html

Delta Lake 的所有命令都可以在 Anaconda shell-prompt 中正常运行。

在 jupyter notebook 上,引用 deltalake 表会出错。这是我在 Jupyter Notebook 上运行的代码 -

df_advisorMetrics.write.mode("overwrite").format("delta").save("/DeltaLake/METRICS_F_DELTA")
spark.sql("create table METRICS_F_DELTA using delta location '/DeltaLake/METRICS_F_DELTA'")

下面是我在笔记本开始时用来连接到 pyspark 的代码 -

import findspark
findspark.init()
findspark.find()

import pyspark
findspark.find()

以下是我得到的错误:

Py4JJavaError:调用 o116.save 时出错。 :java.lang.ClassNotFoundException:找不到数据源:delta。请在http://spark.apache.org/third-party-projects.html查找包

有什么建议吗?

【问题讨论】:

  • 这里有什么建议吗?

标签: pyspark jupyter-notebook delta-lake


【解决方案1】:

我创建了一个 Google Colab/Jupyter Notebook 示例来展示如何运行 Delta Lake。

https://github.com/prasannakumar2012/spark_experiments/blob/master/examples/Delta_Lake.ipynb

它具有运行所需的所有步骤。这使用最新的 spark 和 delta 版本。请相应地更改版本。

【讨论】:

    【解决方案2】:

    一个潜在的解决方案是遵循Import PySpark packages with a regular Jupyter notebook 中提到的技术。

    另一个可能的解决方案是下载 delta-core JAR 并将其放在 $SPARK_HOME/jars 文件夹中,这样当您运行 jupyter notebook 时,它会自动包含 Delta Lake JAR。

    【讨论】:

    • 很遗憾这个提议并没有解决问题
    【解决方案3】:

    我一直在 Jupyter 笔记本上使用 DeltaLake。

    在运行 Python 3.x 的 Jupyter 笔记本中尝试以下操作。

    ### import Spark libraries
    from pyspark.sql import SparkSession
    import pyspark.sql.functions as F
    
    ### spark package maven coordinates - in case you are loading more than just delta
    spark_packages_list = [
        'io.delta:delta-core_2.11:0.6.1',
    ]
    spark_packages = ",".join(spark_packages_list)
    
    ### SparkSession 
    spark = (
        SparkSession.builder
        .config("spark.jars.packages", spark_packages)
        .config("spark.delta.logStore.class", "org.apache.spark.sql.delta.storage.S3SingleDriverLogStore")
        .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
        .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") 
        .getOrCreate()
    )
    
    sc = spark.sparkContext
    
    ### Python library in delta jar. 
    ### Must create sparkSession before import
    from delta.tables import *
    

    假设你有一个 spark 数据框 df

    HDFS

    保存

    ### overwrite, change mode="append" if you prefer
    (df.write.format("delta")
    .save("my_delta_file", mode="overwrite", partitionBy="partition_column_name")
    )
    

    加载

    df_delta = spark.read.format("delta").load("my_delta_file")
    

    AWS S3 对象存储

    初始 S3 设置

    ### Spark S3 access
    hdpConf = sc._jsc.hadoopConfiguration()
    user = os.getenv("USER")
    
    ### Assuming you have your AWS credentials in a jceks keystore.
    hdpConf.set("hadoop.security.credential.provider.path", f"jceks://hdfs/user/{user}/awskeyfile.jceks")
    
    hdpConf.set("fs.s3a.fast.upload", "true")
    
    ### optimize s3 bucket-level parquet column selection
    ### un-comment to use
    # hdpConf.set("fs.s3a.experimental.fadvise", "random")
    
    
    ### Pick one upload buffer option
    hdpConf.set("fs.s3a.fast.upload.buffer", "bytebuffer") # JVM off-heap memory
    # hdpConf.set("fs.s3a.fast.upload.buffer", "array") # JVM on-heap memory
    # hdpConf.set("fs.s3a.fast.upload.buffer", "disk") # DEFAULT - directories listed in fs.s3a.buffer.dir
    
    s3_bucket_path = "s3a://your-bucket-name"
    s3_delta_prefix = "delta"  # or whatever
    

    保存

    ### overwrite, change mode="append" if you prefer
    (df.write.format("delta")
    .save(f"{s3_bucket_path}/{s3_delta_prefix}/", mode="overwrite", partitionBy="partition_column_name")
    )
    

    加载

    df_delta = spark.read.format("delta").load(f"{s3_bucket_path}/{s3_delta_prefix}/")
    

    火花提交

    不直接回答原始问题,但为了完整起见,您也可以执行以下操作。

    将以下内容添加到您的 spark-defaults.conf 文件中

    spark.jars.packages                 io.delta:delta-core_2.11:0.6.1
    spark.delta.logStore.class          org.apache.spark.sql.delta.storage.S3SingleDriverLogStore
    spark.sql.extensions                io.delta.sql.DeltaSparkSessionExtension
    spark.sql.catalog.spark_catalog     org.apache.spark.sql.delta.catalog.DeltaCatalog
    

    参考 spark-submit 命令中的 conf 文件

    spark-submit \
    --properties-file /path/to/your/spark-defaults.conf \
    --name your_spark_delta_app \
    --py-files /path/to/your/supporting_pyspark_files.zip \
    --class Main /path/to/your/pyspark_script.py
    

    【讨论】:

      猜你喜欢
      • 2017-06-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-03-29
      • 2023-03-04
      • 2021-11-23
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多