使用 pyspark 将镶木地板数据写入 csv 时出现“不支持的编码:DELTA_BYTE_ARRAY”

Pri*_*i31 3 apache-spark-sql pyspark

我想将二进制格式的 parquet 文件转换为 csv 文件。我在 Spark 中使用以下命令。

sqlContext.setConf("spark.sql.parquet.binaryAsString","true")

val source =  sqlContext.read.parquet("path to parquet file")

source.coalesce(1).write.format("com.databricks.spark.csv").option("header","true").save("path to csv")
Run Code Online (Sandbox Code Playgroud)

当我在 HDFS 服务器中启动 Spark 并运行这些命令时,这是有效的。当我尝试将相同的 parquet 文件复制到本地系统并启动 pyspark 并运行这些命令时,出现错误。

我能够将二进制作为字符串属性设置为 true,并且能够读取本地 pyspark 中的镶木地板文件。但是当我执行写入 csv 的命令时,出现以下错误。

2018-10-01 14:45:11 警告 ZlibFactory:51 - 无法加载/初始化 native-zlib 库 2018-10-01 14:45:12 错误 Utils:91 - 中止任务 java.lang.UnsupportedOperationException:不支持的编码: DELTA_BYTE_ARRAY 在 org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.initDataReader(VectorizedColumnReader.java:577) 在 org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.readPageV2(VectorizedColumnReader.java:627) )在 org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.access$100(VectorizedColumnReader.java:47) 在 org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader$1.visit(VectorizedColumnReader.java :550)在org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader $1.visit(VectorizedColumnReader.java:536)在org.apache.parquet.column.page.DataPageV2.accept(DataPageV2.java:141)在 org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.readPage(VectorizedColumnReader.java:536) 在 org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.readBatch(VectorizedColumnReader.java:164)在 org.apache.spark.sql.execution.datasources.parquet.VectorizedParquetRecordReader.nextBatch(VectorizedParquetRecordReader.java:263) 在 org.apache.spark.sql.execution.datasources.parquet.VectorizedParquetRecordReader.nextKeyValue(VectorizedParquetRecordReader.java:161)在 org.apache.spark.sql.execution.datasources.RecordReaderIterator.hasNext(RecordReaderIterator.scala:39) 在 org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.hasNext(FileScanRDD.scala:109)在 org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.nextIterator(FileScanRDD.scala:186) 在 org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.hasNext(FileScanRDD.等级:109)

应该如何解决本地计算机中的此错误,就像在 hdfs 中一样?任何解决这个问题的想法都会有很大帮助。谢谢。

Taw*_*kir 5

您可以尝试禁用 VectorizedReader。

spark.conf.set("spark.sql.parquet.enableVectorizedReader", "false")
Run Code Online (Sandbox Code Playgroud)

这不是解决方案,而是一种解决方法。禁用它的后果将是https://jaceklaskowski.gitbooks.io/mastering-spark-sql/spark-sql-vectorized-parquet-reader.html