【问题标题】:How can I read/write data from Azurite using Spark?如何使用 Spark 从 Azurite 读取/写入数据?
【发布时间】:2021-03-11 00:35:52
【问题描述】:

我尝试使用 Spark 从/向Azurite 读取/写入 Parquet 文件,如下所示:

import com.holdenkarau.spark.testing.DatasetSuiteBase
import org.apache.spark.SparkConf
import org.apache.spark.sql.SaveMode
import org.scalatest.WordSpec

class SimpleAzuriteSpec extends WordSpec with DatasetSuiteBase {
  val AzuriteHost = "localhost"
  val AzuritePort = 10000
  val AzuriteAccountName = "devstoreaccount1"
  val AzuriteAccountKey = "Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw=="
  val AzuriteContainer = "container1"
  val AzuriteDirectory = "dir1"
  val AzuritePath = s"wasb://$AzuriteContainer@$AzuriteAccountName.blob.core.windows.net/$AzuriteDirectory/"

  override final def conf: SparkConf = {
    val cfg = super.conf
    val settings =
      Map(
        s"spark.hadoop.fs.azure.storage.emulator.account.name" -> AzuriteAccountName,
        s"spark.hadoop.fs.azure.account.key.${AzuriteAccountName}.blob.core.windows.net" -> AzuriteAccountKey
      )
    settings.foreach { case (k, v) =>
      cfg.set(k, v)
    }
    cfg
  }

  "Spark" must {
    "write to/read from Azurite" in {
      import spark.implicits._
      val xs = List(Rec(1, "Alice"), Rec(2, "Bob"))
      val inputDs = spark.createDataset(xs)

      inputDs.write
        .format("parquet")
        .mode(SaveMode.Overwrite)
        .save(AzuritePath)

      val ds = spark.read
        .format("parquet")
        .load(AzuritePath)
        .as[Rec]

      ds.show(truncate = false)

      val actual = ds.collect().toList.sortBy(_.id)
      assert(actual == xs)
    }
  }
}

case class Rec(id: Int, name: String)
  • 我已经尝试过 Azurite 3.9.0 和 Azurite 2.7.0(都在 Docker 中)。我可以使用az(也可以使用dockerized)向/从Azurite传输文件。

  • 上面的测试在 Docker 主机上运行。 Azurite 可从 Docker 主机访问。

我正在使用 Spark 2.4.5、Hadoop 2.10.0 和此依赖项:

libraryDependencies += "org.apache.hadoop" % "hadoop-azure" % "2.10.0"

使用az 时,此连接字符串有效:

AZURE_STORAGE_CONNECTION_STRING="DefaultEndpointsProtocol=http;AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw==;BlobEndpoint=http://azurite-3.9.0:10000/devstoreaccount1;QueueEndpoint=http://azurite-3.9.0:10001/devstoreaccount1;"

但我不知道如何在 Spark 中进行配置。

我的问题:如何配置主机、端口、凭据等(在路径中或 SparkConf 中)?

【问题讨论】:

  • 你是在同一台机器上运行 spark 吗?
  • @JimXu 都在同一台机器上运行。 Spark 在 Docker 主机上运行(= 不在 Docker 中)。 Azurite 和 az 在 Docker 中运行。 Azurite 端口可访问(使用curl 验证)。但是必须在某个地方配置该端口。

标签: scala azure apache-spark azure-storage azurite


【解决方案1】:

不要使用 Azurite,只需将这些 Jars 添加到您的 Spark Dockerfile:

# Set JARS env
ENV JARS=${SPARK_HOME}/jars/azure-storage-${AZURE_STORAGE_VER}.jar,${SPARK_HOME}/jars/hadoop-azure-${HADOOP_AZURE_VER}.jar,${SPARK_HOME}/jars/jetty-util-ajax-${JETTY_VER}.jar,${SPARK_HOME}/jars/jetty-util-${JETTY_VER}.jar

RUN echo "spark.jars ${JARS}" >> $SPARK_HOME/conf/spark-defaults.conf

设置你的配置:

spark = SparkSession.builder.config(conf=sparkConf).getOrCreate()
spark.sparkContext._jsc.hadoopConfiguration().set(f"fs.azure.account.key.{ os.environ['AZURE_STORAGE_ACCOUNT'] }.blob.core.windows.net", os.environ['AZURE_STORAGE_KEY'])

然后你就可以阅读了:

val df = spark.read.parquet("wasbs://<container-name>@<storage-account-name>.blob.core.windows.net/<directory-name>")

【讨论】:

  • 问题是关于“如何读取/写入数据到 Azurite”。所以像“不要使用蓝晶石”这样的答案并没有真正的帮助。我想使用 Azurite 进行独立的集成测试——独立于任何 Azure 服务/云等。
【解决方案2】:

是的,这是可能的,但是对于 wasb,应该可以通过 127.0.0.1:10000 访问 azurite(因此,如果它在另一台机器上运行,则端口转发会有所帮助),然后指定以下 spark args 作为示例:

./pyspark --conf "spark.hadoop.fs.defaultFS=wasb://container@azurite" --conf "spark.hadoop.fs.azure.storage.emulator.account.name=azurite"

然后默认文件系统将由您的 azurite 实例备份。

【讨论】:

  • 您使用了哪些版本(Azurite、azure-cli、Spark)?能否提供一个完整的 PySpark 示例?
猜你喜欢
  • 1970-01-01
  • 2023-03-31
  • 2017-04-03
  • 2019-03-14
  • 1970-01-01
  • 1970-01-01
  • 2016-06-29
  • 2016-10-12
  • 1970-01-01
相关资源
最近更新 更多