【问题标题】:Avoid write files for empty partitions in Spark Streaming避免在 Spark Streaming 中为空分区写入文件
【发布时间】:2019-05-01 11:54:21
【问题描述】:

我有从 kafka 分区 (one executor per partition) 读取数据的 Spark Streaming 作业。
我需要将转换后的值保存到 HDFS,但需要避免创建空文件。
我尝试使用 isEmpty,但当并非所有分区都为空时,这无济于事。

附:由于性能下降,重新分区不是可接受的解决方案。

【问题讨论】:

  • 你可以使用 Kafka Connect 来代替......这样你就不需要编写代码,你也不会有空文件
  • @cricket_007 这可能适用于文本数据,但不适用于需要处理和多个输出的我的 avro 管道。现在它适用于 LazyOutputFormat
  • @cricket_007 我有 json,而不是 Kafka 中的 avro。我为每条消息在 avro 中构建了三个具有不同内容的输出。在您发表第一条评论后,我阅读了 confluent.io 上的页面,但仍然认为它不能解决我的问题。

标签: apache-spark hadoop apache-kafka spark-streaming spark-streaming-kafka


【解决方案1】:

该代码仅适用于 PairRDD。

文本代码:

  val conf = ssc.sparkContext.hadoopConfiguration
  conf.setClass("mapreduce.output.lazyoutputformat.outputformat",
    classOf[TextOutputFormat[Text, NullWritable]]
    classOf[OutputFormat[Text, NullWritable]])

  kafkaRdd.map(_.value -> NullWritable.get)
    .saveAsNewAPIHadoopFile(basePath,
      classOf[Text],
      classOf[NullWritable],
      classOf[LazyOutputFormat[Text, NullWritable]],
      conf)

avro 代码:

  val avro: RDD[(AvroKey[MyEvent], NullWritable)]) = ....
  val conf = ssc.sparkContext.hadoopConfiguration

  conf.set("avro.schema.output.key", MyEvent.SCHEMA$.toString)
  conf.setClass("mapreduce.output.lazyoutputformat.outputformat",
    classOf[AvroKeyOutputFormat[MyEvent]],
    classOf[OutputFormat[AvroKey[MyEvent], NullWritable]])

  avro.saveAsNewAPIHadoopFile(basePath,
    classOf[AvroKey[MyEvent]],
    classOf[NullWritable],
    classOf[LazyOutputFormat[AvroKey[MyEvent], NullWritable]],
    conf)

【讨论】:

    猜你喜欢
    • 2018-11-15
    • 2020-05-28
    • 1970-01-01
    • 1970-01-01
    • 2015-11-25
    • 1970-01-01
    • 1970-01-01
    • 2016-01-04
    • 2018-03-21
    相关资源
    最近更新 更多