【问题标题】:How to write to Kafka from Spark with a changed schema without getting exceptions?如何使用更改的架构从 Spark 写入 Kafka 而不会出现异常?
【发布时间】:2018-11-23 23:00:24
【问题描述】:

我正在将 parquet 文件从 Databricks 加载到 Spark:

val dataset = context.session.read().parquet(parquetPath)

然后我执行一些这样的转换:

val df = dataset.withColumn(
            columnName, concat_ws("",
            col(data.columnName), lit(textToAppend)))

当我尝试将其作为 JSON 保存到 Kafka 时(不回到镶木地板!):

df = df.select(
            lit("databricks").alias("source"),
            struct("*").alias("data"))

val server = "kafka.dev.server" // some url
df = dataset.selectExpr("to_json(struct(*)) AS value")
df.write()
        .format("kafka")
        .option("kafka.bootstrap.servers", server)
        .option("topic", topic)
        .save()

我得到以下异常:

org.apache.spark.sql.execution.QueryExecutionException: Parquet column cannot be converted in file dbfs:/mnt/warehouse/part-00001-tid-4198727867000085490-1e0230e7-7ebc-4e79-9985-0a131bdabee2-4-c000.snappy.parquet. Column: [item_group_id], Expected: StringType, Found: INT32
    at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1$$anonfun$prepareNextFile$1.apply(FileScanRDD.scala:310)
    at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1$$anonfun$prepareNextFile$1.apply(FileScanRDD.scala:287)
    at scala.concurrent.impl.Future$PromiseCompletingRunnable.liftedTree1$1(Future.scala:24)
    at scala.concurrent.impl.Future$PromiseCompletingRunnable.run(Future.scala:24)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)
Caused by: org.apache.spark.sql.execution.datasources.SchemaColumnConvertNotSupportedException
    at com.databricks.sql.io.parquet.NativeColumnReader.readBatch(NativeColumnReader.java:448)
    at com.databricks.sql.io.parquet.DatabricksVectorizedParquetRecordReader.nextBatch(DatabricksVectorizedParquetRecordReader.java:330)
    at org.apache.spark.sql.execution.datasources.parquet.VectorizedParquetRecordReader.nextKeyValue(VectorizedParquetRecordReader.java:167)
    at org.apache.spark.sql.execution.datasources.RecordReaderIterator.hasNext(RecordReaderIterator.scala:40)
    at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1$$anonfun$prepareNextFile$1.apply(FileScanRDD.scala:299)
    at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1$$anonfun$prepareNextFile$1.apply(FileScanRDD.scala:287)
    at scala.concurrent.impl.Future$PromiseCompletingRunnable.liftedTree1$1(Future.scala:24)
    at scala.concurrent.impl.Future$PromiseCompletingRunnable.run(Future.scala:24)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

只有在我尝试读取多个分区时才会发生这种情况。例如,在/mnt/warehouse/ 目录中,我有很多镶木地板文件,每个文件都代表来自datestamp 的数据。如果我只阅读其中一个,我不会收到异常,但如果我阅读整个目录,则会发生此异常。

当我进行转换时,我会得到这个,就像上面我更改列的数据类型一样。我怎样才能解决这个问题?我不是想写回 parquet,而是将所有文件从同一源架构转换为新架构并将它们写入 Kafka。

【问题讨论】:

  • 如果你想写一个 kafka 主题,一个解决方案可能是使用生产者

标签: scala apache-spark apache-kafka parquet databricks


【解决方案1】:

你可以在这个link找到说明

它向您展示了将数据写入 kafka 主题的不同方式。

【讨论】:

    【解决方案2】:

    parquet 文件似乎存在问题。文件中的item_group_id 列并非都是相同的数据类型,一些文件将该列存储为字符串,而其他文件则存储为整数。从异常SchemaColumnConvertNotSupportedException的源码我们看到描述:

    parquet reader 发现列类型不匹配时引发异常。

    在github 上的 Spark 测试中可以找到一种简单的方法来复制问题:

    Seq(("bcd", 2)).toDF("a", "b").coalesce(1).write.mode("overwrite").parquet(s"$path/parquet")
    Seq((1, "abc")).toDF("a", "b").coalesce(1).write.mode("append").parquet(s"$path/parquet")
    
    spark.read.parquet(s"$path/parquet").collect()
    

    当然,这只会在一次读取多个文件时发生,或者在上面的测试中附加了更多数据。如果读取单个文件,则不会出现列的数据类型之间的不匹配问题。


    解决问题的最简单的方法是在写入文件时确保所有文件的列类型正确。

    替代方法是分别读取所有 parquet 文件,更改模式以匹配,然后将它们与union 结合。一个简单的方法是调整架构:

    // Specify the files and read as separate dataframes
    val files = Seq(...)
    val dfs = files.map(file => spark.read.parquet(file))
    
    // Specify the schema (here the schema of the first file is used)
    val schema = dfs.head.schema
    
    // Create new columns with the correct names and types
    val newCols = schema.map(c => col(c.name).cast(c.dataType))
    
    // Select the new columns and merge the dataframes
    val df = dfs.map(_.select(newCols: _*)).reduce(_ union _)
    

    【讨论】:

    • 非常感谢,这就是问题所在。源文件有错误的数据。感谢您指出这一点!
    猜你喜欢
    • 2019-02-17
    • 2020-03-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-08-19
    • 2019-09-20
    • 2021-09-18
    • 2017-02-27
    相关资源
    最近更新 更多