【问题标题】:Flink convert to parquet errorFlink 转换为 parquet 错误
【发布时间】:2016-09-09 13:46:59
【问题描述】:

我正在尝试使用 flink 将 csv 文件编写为镶木地板。 我正在使用以下代码并得到错误。

val parquetFormat = new HadoopOutputFormat[Void, String](new AvroParquetOutputFormat, job)
FileOutputFormat.setOutputPath(job, new Path(outputPath))

我收到以下构建错误。有人可以帮忙吗?

类型不匹配;发现:parquet.avro.AvroParquetOutputFormat 必需的: org.apache.hadoop.mapreduce.OutputFormat[Void,String] ingestion.scala /flink-scala/src/main/scala/com/sc/edl/flink 行 75 斯卡拉问题

【问题讨论】:

    标签: hadoop apache-flink parquet


    【解决方案1】:

    您想创建一个需要OutputFormat[Void, String]HadoopOutputFormat[Void, String]

    您提供一个扩展ParquetOutputFormat<IndexedRecord>AvroParquetOutputFormatParquetOutputFormat 定义为ParquetOutputFormat<T> extends FileOutputFormat<Void, T>

    因此,您提供OutputFormat[Void, IndexedRecord],而HadoopOutputFormat[Void, String] 需要OutputFormat[Void, String]

    你应该把parquetFormat改成

    val parquetFormat = new HadoopOutputFormat[Void, IndexedRecord](
      new AvroParquetOutputFormat, job)
    FileOutputFormat.setOutputPath(job, new Path(outputPath))
    

    如果您要写出的DataSet 不是(Void, IndexedRecord) 类型,您应该添加一个MapFunction,将您的数据转换成(Void, IndexedRecord) 对。

    【讨论】:

    • 谢谢 Fabian,对不起,我是新手,请您提供正确的语法或有什么问题
    • 我扩展了我的答案
    • Niki 你能解决吗。我也有同样的问题,你能帮帮我吗?
    • @FabianHueske 您好 Fabian,感谢您的回答。你能给我一个关于如何转换为 (Void,IndexedRecord) 的提示吗?我有 DataSet[MyType] 并使用 MapFunction,但我无法实例化 Void。谢谢
    【解决方案2】:

    问题仍然存在,因为 Flink Tuple 目前不支持 NULL Keys。 会出现以下错误: Caused by: org.apache.flink.types.NullFieldException: Field 1 is null, but expected to hold a value.

    更好的选择是使用 KiteSDK,如本示例中所述: https://github.com/nezihyigitbasi/FlinkParquet 因此,如果您需要动态模式,那么这种方法将不起作用,因为您需要严格遵守模式。此外,这更适合阅读而不是写作。

    Spark DataFrame 与 Parquet 配合得非常好,不仅在 API 方面,而且在性能方面。但是如果要使用 Flink,那么你需要等待 flink 社区发布 api 或者编辑自己的 parquet-hadoop 代码,这可能是一个很大的努力。

    目前仅实现了这些连接器 https://github.com/apache/flink/tree/master/flink-connectors 所以,我个人的建议是,如果你可以使用 spark,那就去吧,考虑到生产用例,它有更成熟的 api。当您坚持使用 flink 的基本需求时,您可能还会卡在其他地方。

    到目前为止,不要浪费时间寻找 Flink 的解决方法,我已经浪费了很多关键时间,而不是使用 Hive、Spark 或 MR 等标准选项。

    【讨论】:

      【解决方案3】:

      扩展另一个答案,您可以通过下降到 Java 来实例化所需的 Void 类型:

      // in src/main/java/com/yourOrg/FlinkUtils.java
      public class FlinkUtils {
          /* Stupid hack because we can't instantiate Void in Scala */
          public static Void getVoid() {
              return null;
          }
      }
      
      // src/main/scala/com/yourOrg/FlinkJob.scala
      
      val voidKeyedDataset = ds.map((FlinkUtils.getVoid, _))
      voidKeyedDataset.output(...)
      

      【讨论】:

        猜你喜欢
        • 2019-09-07
        • 2022-06-13
        • 2016-12-06
        • 2018-12-15
        • 1970-01-01
        • 1970-01-01
        • 2019-05-15
        • 2014-06-22
        相关资源
        最近更新 更多