首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >在SparkSQL中使用Avro模式和Parquet格式进行读写

在SparkSQL中使用Avro模式和Parquet格式进行读写
EN

Stack Overflow用户
提问于 2017-01-04 03:13:04
回答 1查看 2.2K关注 0票数 5

我正在尝试从SparkSQL中写入和读取镶木面板文件。出于模式演变的原因,我希望在写入和读取时使用Avro模式。

我的理解是,这在Spark之外(或在Spark内手动)是可能的,例如使用AvroParquetWriter和Avro的通用API。但是,我想使用SparkSQL的write()和read()方法(它们与DataFrameWriter和DataFrameReader一起工作),它们与SparkSQL很好地集成在一起(我将编写和读取Dataset)。

我无论如何也想不出如何做到这一点,我想知道这是否可能。SparkSQL拼接格式似乎只支持“压缩”和"mergeSchema“选项--也就是说,没有用于指定替代模式格式或替代模式的选项。换句话说,似乎没有办法通过Avro模式使用SparkSQL应用编程接口来读/写拼图文件。但也许我只是错过了什么?

为了澄清,我也理解,这将基本上只是添加Avro模式到拼花的元数据写入,并将添加一个更多的翻译层读取(Parquet格式-> Avro模式-> SparkSQL内部格式),但将特别允许我添加缺失列的默认值( Avro模式支持,但Parquet模式不)。

此外,我并不是在寻找一种方法来转换Avro到拼图,或者拼图到Avro (而是一种同时使用它们的方式),我也不是在寻找一种在SparkSQL中读/写普通Avro的方法(你可以使用databricks/spark-avro来做到这一点)。

EN

回答 1

Stack Overflow用户

发布于 2017-01-06 02:20:10

我也在做类似的事情。我使用avro模式来写入拼图文件,但是,不要将其读取为avro。但同样的技术也应该适用于read。我不确定这是否是最好的方法,但不管怎样:我有一个具有avro模式的AvroData.avsc。

代码语言:javascript
复制
KafkaUtils.createDirectStream[String,Array[Byte],StringDecoder,DefaultDecoder,Tuple2[String, Array[Byte]]](ssc, kafkaProps, fromOffsets, messageHandler)


kafkaArr.foreachRDD  { (rdd,time) 
       => { val schema =  SchemaConverters.toSqlType(AvroData.getClassSchema).dataType.asInstanceOf[StructType] val ardd = rdd.mapPartitions{itr =>
              itr.map { r =>
try {
                    val cr = avroToListWithAudit(r._2, offsetSaved, loadDate, timeNow.toString)
                    Row.fromSeq(cr.toArray)
    } catch{
      case e:Exception => LogHandler.log.error("Exception while converting to Avro" + e.printStackTrace())
      System.exit(-1)
      Row(0)  //This is just to allow compiler to accept. On exception, the application will exit before this point
} 
} 
}


  public static List avroToListWithAudit(byte[] kfkBytes, String kfkOffset, String loaddate, String loadtime ) throws IOException {
        AvroData av = getAvroData(kfkBytes);
        av.setLoaddate(loaddate);
        av.setLoadtime(loadtime);
        av.setKafkaOffset(kfkOffset);
        return avroToList(av);
    }



 public static List avroToList(AvroData a) throws UnsupportedEncodingException{
        List<Object> l = new ArrayList<>();
        for (Schema.Field f : a.getSchema().getFields()) {
            String field = f.name().toString();
            Object value = a.get(f.name());
            if (value == null) {
                //System.out.println("Adding null");
                l.add(""); 
            }
            else {
                switch (f.schema().getType().getName()){
                    case "union"://System.out.println("Adding union");
                        l.add(value.toString());
                        break;

                    default:l.add(value);
                        break;
                }

            }
        }
        return l;
    }

getAvroData方法需要有从原始字节构造avro对象的代码。我还试图找出一种方法来做到这一点,而不必显式地指定每个属性设置器,但似乎没有这样的方法。

代码语言:javascript
复制
public static AvroData getAvroData (bytes)
{
AvroData av = AvroData.newBuilder().build();
        try {
            av.setAttr(String.valueOf("xyz"));
        .....
    }
   } 

希望能有所帮助

票数 0
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/41450739

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档