【问题标题】:How to read binary avro fileData, with Source in akka?如何使用akka中的Source读取二进制avro fileData?
【发布时间】:2019-06-18 20:23:24
【问题描述】:

我正在尝试使用来自 akka Streams 的 Source 读取 avro 文件。

akka 流中的源读取类似 FileIO.FromPath(File) 的数据,它将根据 (\n) 字符读取和分隔行,而 avro 是如何工作的?

流程:

    object AvroFlow  {
    def apply(jobDate: String): Flow[GenericRecord, GenericRecord, NotUsed] = {
            Flow[GenericRecord].map { rec => rec.put("date", "20190812") rec}       
    }
  }

图表:

object AvroRunner {
    def build (src: Source[GenericRecord, NotUsed],
                                     flw: Flow[GenericRecord, GenericRecord, NotUsed],
                                     snk:Flow[GenericRecord, Future[Done])
    : AvroRunner = {
      new AvroRunner(srtc,flw,snk)
    }
  }
class AvroRunner private(src: Source[GenericRecord, NotUsed],
                                     flw: Flow[GenericRecord, GenericRecord, NotUsed],
                                     snk:Flow[GenericRecord, Future[Done]){
  import scala.concurrent.ExecutionContext.Implicits.global
  val GraphRunner = RunnableGraph.fromGraph(GraphDSL.create() {implicit builder =>
    import GraphDSL.Implicits._
    src ~> flw ~> snk
    ClosedShape
  })
}

【问题讨论】:

    标签: scala io akka avro akka-stream


    【解决方案1】:

    创建 avro 数据对象的 akka Source 的最简单方法不是来自原始二进制文件本身。相反,从 avro 库提供的DataFileReader 创建源。

    the documentation 我们首先从java.io.File 生成器创建文件阅读器:

    def createFileReader[T : ClassTag](fileGenerator : () => File) : DataFileReader[T] = 
      new DataFileReader[T](file(), new SpecificDatumReader[T](classTag[T].runtimeClass))
    

    这可以用来创建一个 scala Iterator:

    def dataFileReaderToIterator[T](dataFileReader : DataFileReader[T]) : Iterator[T] = 
      new Iterator[T] {
        override def hasNext : Boolean = dataFileReader.hasNext
    
        override def next() : T = dataFileReader.next
      }
    

    我们现在可以从文件生成器构造一个流 Source:

    def fileToAvroSource[T](fileGenerator : () => File) : Source[T, _] = 
      Source.fromIterator[T](() => dataFileReaderToIterator[T](createFileReader(fileGenerator)))
    

    背压?

    似乎 avro 正在使用标准的 BufferedReader/OutputStream 技术来读取 File。因此,上述实现应该提供一直到文件源的背压。但是,我还没有确认是这种情况......

    【讨论】:

    • 嗨拉蒙,感谢您的回复。我从Iterator() 创建了 Source,但是当我将 Source 连接到 Flow[GenericRecord,GenericRecord,_] 时,Flow 没有得到任何记录。但是在 Source 中,我能够打印记录。我在这里遗漏的任何东西......例如:Flow[GenericRecord].map { x.put(dt,"20190624")}
    • @Aravinda 欢迎您。我必须查看您的 Flow 代码和流代码才能调试...
    • @J Romero Vigil,我已经在上面的标签中包含了代码以及问题。有什么我想念的吗?
    猜你喜欢
    • 2015-03-06
    • 1970-01-01
    • 1970-01-01
    • 2017-04-26
    • 2014-03-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多