【问题标题】:Binary File read in Google Dataflow在 Google Dataflow 中读取的二进制文件
【发布时间】:2018-10-11 14:13:05
【问题描述】:

我需要在谷歌数据流中读取二进制文件, 我只需要读取文件并将每 64 字节解析为一条记录,并在数据流中每 64 字节二进制文件的每个字节中应用一些逻辑。

我在 spark 中尝试过同样的事情,代码 smape 如下:

 def main(args: Array[String]): Unit = {

    val spark = SparkSession
      .builder()
      .appName("RecordSplit")
      .master("local[*]")
      .getOrCreate()

    val df = spark.sparkContext.binaryRecords("< binary-file-path>", 64)

    val Table = df.map(rec => {
      val c1= (convertHexToString(rec(0)))
      val c2= convertBinaryToInt16(rec, 48)
      val c3= rec(59)
      val c4= convertHexToString(rec(50)) match {
        case str =>
          if (str.startsWith("c"))
            2020 + str.substring(1).toInt
          else if (str.startsWith("b"))
            2010 + str.substring(1).toInt
          else if (str.startsWith("b"))
            2000 + str.substring(1).toInt
        case _ => 1920
      }

【问题讨论】:

  • 欢迎来到 SO。请提供一个最小、完整和可验证的示例。 向我们展示您最近尝试的代码以及您遇到的问题。并解释为什么结果不是你所期望的。编辑您的问题以包含代码,请不要在评论中添加它,因为它可能不可读。 stackoverflow.com/help/mcve 最好展示实际发生的事情,而不是描述您期望发生的事情。请包含代码和输出作为您问题的内容,而不是图片或外部链接。

标签: google-cloud-platform google-cloud-dataflow apache-beam dataflow


【解决方案1】:

我会推荐以下内容:

  • 如果您不限于 python/scala,OffsetBasedSource(FileBasedSource 是一个子类)可以满足您的需求,因为它使用偏移量来定义开始和结束位置。

  • TikaIO 可以处理元数据,但它可以根据文档读取二进制数据。

  • 示例dataflow-opinion-analysis 包含要从任意字节位置读取的信息。

  • 还有其他文档可以创建自定义 Read implementation。您可能需要考虑查看这些Beam examples,以获取有关如何实现自定义源的指导,例如python example。

另一种方法是在管道外部(内存中)创建 64 字节数组,然后创建 PCollection from memory,请记住文档建议将其用于单元测试。

【讨论】:

  • 您也可以作为 SplittableDoFn 执行此操作,或者(假设您有足够的文件来获得所需的并行度)使用 FileIO.match() 后跟一个读取发出 64 字节块的文件的 DoFn . (后一种解决方案不会提供与使用 OffsetBasedSource 一样好的拆分行为,但可能会更快。)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-12-26
  • 2021-12-08
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多