【问题标题】:Upgrading from Flink 1.3.2 to 1.4.0 hadoop FileSystem and Path issues从 Flink 1.3.2 升级到 1.4.0 hadoop 文件系统和路径问题
【发布时间】:2017-12-27 17:43:37
【问题描述】:

我最近尝试从 Flink 1.3.2 升级到 1.4.0,但我遇到了一些问题,无法再导入 org.apache.hadoop.fs.{FileSystem, Path}。问题出现在两个地方:

ParquetWriter:

import org.apache.avro.Schema
import org.apache.avro.generic.GenericRecord
import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.flink.streaming.connectors.fs.Writer
import org.apache.parquet.avro.AvroParquetWriter
import org.apache.parquet.hadoop.ParquetWriter
import org.apache.parquet.hadoop.metadata.CompressionCodecName

class AvroWriter[T <: GenericRecord]() extends Writer[T] {

  @transient private var writer: ParquetWriter[T] = _
  @transient private var schema: Schema = _

  override def write(element: T): Unit = {
    schema = element.getSchema
    writer.write(element)
  }

  override def duplicate(): AvroWriter[T] = new AvroWriter[T]()

  override def close(): Unit = writer.close()

  override def getPos: Long = writer.getDataSize

  override def flush(): Long = writer.getDataSize

  override def open(fs: FileSystem, path: Path): Unit = {
    writer = AvroParquetWriter.builder[T](path)
      .withSchema(schema)
      .withCompressionCodec(CompressionCodecName.SNAPPY)
      .build()
  }

}

自定义桶:

import org.apache.flink.streaming.connectors.fs.bucketing.Bucketer
import org.apache.flink.streaming.connectors.fs.Clock
import org.apache.hadoop.fs.{FileSystem, Path}
import java.io.ObjectInputStream
import java.text.SimpleDateFormat
import java.util.Date

import org.apache.avro.generic.GenericRecord

import scala.reflect.ClassTag

class RecordFieldBucketer[T <: GenericRecord: ClassTag](dateField: String = null, dateFieldFormat: String = null, bucketOrder: Seq[String]) extends Bucketer[T] {

  @transient var dateFormatter: SimpleDateFormat = _

  private def readObject(in: ObjectInputStream): Unit = {
    in.defaultReadObject()
    if (dateField != null && dateFieldFormat != null) {
      dateFormatter = new SimpleDateFormat(dateFieldFormat)
    }
  }

  override def getBucketPath(clock: Clock, basePath: Path, element: T): Path = {
    val partitions = bucketOrder.map(field => {
      if (field == dateField) {
        field + "=" + dateFormatter.format(new Date(element.get(field).asInstanceOf[Long]))
      } else {
        field + "=" + element.get(field)
      }
    }).mkString("/")
    new Path(basePath + "/" + partitions)
  }

}

我注意到 Flink 现在有:

import org.apache.flink.core.fs.{FileSystem, Path}

但新的Path 似乎不适用于AvroParquetWritergetBucketPath 方法。我知道 Flink 的 FileSystem 和 Hadoop 依赖项发生了一些变化,我只是不确定我需要导入什么才能让我的代码再次工作。

我什至需要使用 Hadoop 依赖项,还是现在有不同的方式将 Parquet 文件写入和存储到 s3?

build.sbt:

val flinkVersion = "1.4.0"

libraryDependencies ++= Seq(
  "org.apache.flink" %% "flink-scala" % flinkVersion % Provided,
  "org.apache.flink" %% "flink-streaming-scala" % flinkVersion % Provided,
  "org.apache.flink" %% "flink-connector-kafka-0.10" % flinkVersion,
  "org.apache.flink" %% "flink-connector-filesystem" % flinkVersion,
  "org.apache.flink" % "flink-metrics-core" % flinkVersion,
  "org.apache.flink" % "flink-metrics-graphite" % flinkVersion,
  "org.apache.kafka" %% "kafka" % "0.10.0.1",
  "org.apache.avro" % "avro" % "1.7.7",
  "org.apache.parquet" % "parquet-hadoop" % "1.8.1",
  "org.apache.parquet" % "parquet-avro" % "1.8.1",
  "io.confluent" % "kafka-avro-serializer" % "3.2.2",
  "com.fasterxml.jackson.core" % "jackson-core" % "2.9.2"
)

【问题讨论】:

    标签: apache-flink avro parquet flink-streaming


    【解决方案1】:

    构建“Hadoop-Free-Flink”是 1.4 版本的一项主要功能。 您所要做的就是将 hadoop 依赖项包含到您的类路径中或引用 changelogs:

    ...这也意味着,如果您使用了 HDFS 的连接器,例如 BucketingSink 或 RollingSink,您现在必须确保使用带有捆绑 Hadoop 依赖项的 Flink 发行版,或者确保在以下情况下包含 Hadoop 依赖项为您的应用程序构建一个 jar 文件。

    【讨论】:

    • 好吧,我会尝试追踪我需要包含的依赖项——除非你知道。我还想知道是否我什至需要包含将 parquet 写入 s3 的依赖项,或者现在在 Flink 1.4 中是否有其他方法可以做到这一点?
    【解决方案2】:

    hadoop-commons 项目中可以找到必要的org.apache.hadoop.fs.{FileSystem, Path} 类。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-12-04
      • 2019-09-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-09-30
      • 2017-06-16
      相关资源
      最近更新 更多