【问题标题】:Spark - write Avro fileSpark - 编写 Avro 文件
【发布时间】:2015-11-23 18:53:15
【问题描述】:

在这样的流程中使用 Spark(使用 Scala API)编写 Avro 文件的常见做法是什么:

  1. 从 HDFS 解析一些日志文件
  2. 为每个日志文件应用一些业务逻辑并生成 Avro 文件(或者可能合并多个文件)
  3. 将 Avro 文件写入 HDFS

我尝试使用 spark-avro,但没有多大帮助。

val someLogs = sc.textFile(inputPath)

val rowRDD = someLogs.map { line =>
  createRow(...)
}

val sqlContext = new SQLContext(sc)
val dataFrame = sqlContext.createDataFrame(rowRDD, schema)
dataFrame.write.avro(outputPath)

这失败并出现错误:

org.apache.spark.sql.AnalysisException: 
      Reference 'StringField' is ambiguous, could be: StringField#0, StringField#1, StringField#2, StringField#3, ...

【问题讨论】:

  • 您能说得更具体些吗?例如为什么 `spark-avro` 对你不起作用?
  • 我没有成功使用 Avro 生成的 java 代码和 spark-avro。此外,当我使用 Schema API 时,我会收到此类错误:org.apache.spark.sql.AnalysisException: Reference 'StringField' is ambiguous, could be: StringField#0, StringField#1, StringField#2, StringField#3 ,
  • @d4rkang3l 你确定问题出在 avro 序列化上吗?生成的dataFrame是否没有问题?

标签: apache-spark avro


【解决方案1】:

Databricks 提供了库 spark-avro,它可以帮助我们读写 Avro 数据。

dataframe.write.format("com.databricks.spark.avro").save(outputPath)

【讨论】:

    【解决方案2】:

    Spark 2 和 Scala 2.11

    import com.databricks.spark.avro._
    import org.apache.spark.sql.SparkSession
    
    val spark = SparkSession.builder().master("local").getOrCreate()
    
    // Do all your operations and save it on your Dataframe say (dataFrame)
    
    dataFrame.write.avro("/tmp/output")
    

    Maven 依赖项

    <dependency>
        <groupId>com.databricks</groupId>
        <artifactId>spark-avro_2.11</artifactId>
        <version>4.0.0</version> 
    </dependency>
    

    【讨论】:

    • 在 Java 8 中:df2.write().format("com.databricks.spark.avro").mode(SaveMode.Overwrite).save(f.getOutputPath());
    • 我怀疑 .avro("/tmp/output") 是否有效。我认为应该使用这种方式:write.format("avro").save(/output/path)
    【解决方案3】:

    您需要启动 spark shell 以包含 avro 包。推荐用于较低版本

    $SPARK_HOME/bin/spark-shell --packages com.databricks:spark-avro_2.11:4.0.0

    然后用 df 写成 avro 文件-

    dataframe.write.format("com.databricks.spark.avro").save(outputPath)
    

    并在 hive 中写为 avro 表 -

    dataframe.write.format("com.databricks.spark.avro").saveAsTable(hivedb.hivetable_avro)
    

    【讨论】:

    • Paul “追加”到表的语法,我似乎无法正确理解。如果可能的话,你能验证一下吗?
    • '.saveAsTable(hivedb.hivetable_avro)' 部分将出现以下警告问题:WARN hive.HiveExternalCatalog:找不到数据源提供程序 com.databricks.spark.avro 的相应 Hive SerDe。以 Spark SQL 特定格式将数据源表 dbName.Suchas.hivedb.tableName.suchas.hivetable_avro 持久化到 Hive 元存储中,这与 Hive 不兼容。至少 spark2 发生了这种情况
    猜你喜欢
    • 1970-01-01
    • 2014-01-03
    • 1970-01-01
    • 2021-07-01
    • 2021-05-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-05-11
    相关资源
    最近更新 更多